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::ops::Deref;
46use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering as Atomic};
47use std::sync::{Arc, Mutex};
48
49use rudb_common::{Error, LogicalType, Result};
50use rudb_metrics::{LoadProfile, Stage};
51use rudb_storage::Range;
52use rudb_vector::{Bitmap, Chunk, Data, StringColumn, Validity, Vector};
53
54use super::{
55    ColumnStripe, DICTIONARY_CHECK_SEED, DICTIONARY_DECIDE_ROWS, DICTIONARY_DISTINCT_IN_TEN,
56    EncodedBlock, GlobalDictionary, MAX_ENCODE_WORKERS, MAX_PAGE, Part, PendingChunk, STRIPE_PARTS,
57    Settling, Spread, Unencoded, Writer, checksum, coded_page, invalid, push_validity,
58    seeded_checksum, stats, unique_codes, weight,
59};
60
61/// How many stripes are being prepared or paged right now, across every writer in the process.
62///
63/// A stripe's columns are spread over threads of their own, which is what a writer being fed by
64/// one caller needs, because that caller is the only one encoding. Thirty two callers each doing
65/// that at once would be a thousand threads on a machine with thirty two cores. So each one takes
66/// its share of the machine: the cores over however many stripes are being worked on right now.
67static BUSY: AtomicUsize = AtomicUsize::new(0);
68
69/// What all of a table's global dictionaries may hold at once before the fastest growing one is
70/// demoted, from section 5.5 of the encoding spec.
71///
72/// Dictionaries are the one thing a load holds that grows with the table rather than with the
73/// stripe. `hits` has tens of millions of distinct `URL`, `Title`, `Referer` and `SearchPhrase`
74/// values, at forty five to seventy bytes each for the lookup alone, which is past the whole two
75/// gigabyte bound on its own.
76pub const DICTIONARY_CAP_BYTES: u64 = 512 * 1024 * 1024;
77
78/// Which varchar columns still code against a global dictionary, and what their dictionaries hold
79/// between them.
80///
81/// Shared by a writer with every [`Preparer`] and [`Merger`] it hands out. The flags are what a
82/// stripe prepared later reads to decide whether to code a column at all, and the rest is what a
83/// merge reads to decide whether its column should stop, see [`demotes`].
84#[derive(Debug)]
85pub(crate) struct Coding {
86    flags: Box<[AtomicBool]>,
87    /// What each column's dictionary grew by in the last stripe merged into it.
88    growth: Box<[AtomicU64]>,
89    /// What every dictionary held the last time it was merged into, added up.
90    held: AtomicU64,
91    cap: AtomicU64,
92    /// The lowest distinct count ceiling a stripe of each column has reported, `u64::MAX` until one
93    /// has.
94    ///
95    /// Every stripe prepared for a writer is merged into the writer's statistics, so a stripe's full
96    /// sketch is one of the parts of the union, and no hash above its ceiling can be in the union's
97    /// bottom k. A stripe started later drops those hashes rather than hashing them into a table
98    /// that the union would trim them from anyway. Without it, every stripe filled a sketch from
99    /// empty, and `Sketch::insert` was about 2% of a `lineitem` load's samples.
100    ceilings: Box<[AtomicU64]>,
101}
102
103impl Coding {
104    pub(crate) fn new(flags: impl IntoIterator<Item = bool>) -> Self {
105        let flags = flags.into_iter().map(AtomicBool::new).collect::<Box<[_]>>();
106        let growth = flags.iter().map(|_| AtomicU64::new(0)).collect();
107        let ceilings = flags.iter().map(|_| AtomicU64::new(u64::MAX)).collect();
108        Self {
109            flags,
110            growth,
111            held: AtomicU64::new(0),
112            cap: AtomicU64::new(DICTIONARY_CAP_BYTES),
113            ceilings,
114        }
115    }
116
117    /// Sets what the dictionaries may hold between them before one is demoted.
118    pub(crate) fn cap(&self, bytes: u64) {
119        self.cap.store(bytes, Atomic::Relaxed);
120    }
121
122    /// Moves what one dictionary is counted for from `before` to `now`, and hands back the new
123    /// total.
124    fn recount(&self, before: u64, now: u64) -> u64 {
125        if now >= before {
126            self.held.fetch_add(now - before, Atomic::Relaxed) + (now - before)
127        } else {
128            self.held.fetch_sub(before - now, Atomic::Relaxed).saturating_sub(before - now)
129        }
130    }
131
132    /// Whether the column grew the most in the stripes last merged into each column.
133    ///
134    /// Read without a lock across columns that may be merging at the same moment, so it is about
135    /// the last stripe or the one before it. That is close enough for choosing which column to
136    /// stop: a column that grows fastest keeps doing so, and it is asked again next stripe. A column
137    /// that did not grow at all is never the one, since stopping it would free nothing.
138    fn grew_most(&self, index: usize) -> bool {
139        let mine = self.growth[index].load(Atomic::Relaxed);
140        mine > 0 && self.growth.iter().all(|other| other.load(Atomic::Relaxed) <= mine)
141    }
142}
143
144impl Deref for Coding {
145    type Target = [AtomicBool];
146
147    fn deref(&self) -> &[AtomicBool] {
148        &self.flags
149    }
150}
151
152/// One stripe's share of the machine, held for as long as the stripe is being worked on.
153struct Share(usize);
154
155impl Share {
156    fn take(columns: usize, parts: usize) -> Self {
157        let busy = BUSY.fetch_add(1, Atomic::Relaxed) + 1;
158        let cores =
159            std::thread::available_parallelism().map_or(1, usize::from).min(MAX_ENCODE_WORKERS);
160        // A stripe of one part is one small page a column, which is less than a thread is worth.
161        let workers = if parts <= 1 { 1 } else { (cores / busy).clamp(1, columns.max(1)) };
162        Self(workers)
163    }
164}
165
166impl Drop for Share {
167    fn drop(&mut self) {
168        BUSY.fetch_sub(1, Atomic::Relaxed);
169    }
170}
171
172/// Encodes stripes for one [`Writer`] without the writer.
173///
174/// Handed out by [`Writer::preparer`] and cheap to hold. It shares with its writer which varchar
175/// columns still have a global dictionary, so a stripe prepared after the first one decided a
176/// column should not have one is encoded plainly from the start.
177#[derive(Debug, Clone)]
178pub struct Preparer {
179    types: Vec<LogicalType>,
180    coded: Arc<Coding>,
181    profile: Option<Arc<LoadProfile>>,
182}
183
184/// A stripe that has been through [`Preparer::prepare`] and is waiting for [`Writer::merge`].
185#[derive(Debug)]
186pub struct Prepared {
187    parts: Vec<Part>,
188    types: Vec<LogicalType>,
189    columns: Vec<Column>,
190    gathers: Vec<Option<stats::Gather>>,
191    profile: Option<Arc<LoadProfile>>,
192}
193
194/// A stripe that has been through [`Writer::merge`] and is waiting for [`Merged::pages`].
195#[derive(Debug)]
196pub struct Merged {
197    parts: Vec<Part>,
198    columns: Vec<Merge>,
199    blocks: Vec<Unencoded>,
200    profile: Option<Arc<LoadProfile>>,
201    /// Whether the rows are counted into the table yet. [`Writer::merge`] counts them, and a
202    /// [`Merger`] leaves them for [`Writer::write`], since it has no table to count them into.
203    counted: bool,
204}
205
206/// A stripe that has been through [`Merged::pages`] and is waiting for [`Writer::write`].
207#[derive(Debug)]
208pub struct Paged {
209    parts: Vec<Part>,
210    columns: Vec<ColumnStripe>,
211    /// Encoded dictionary blocks, each with its column and block number.
212    blocks: Vec<(usize, usize, EncodedBlock)>,
213    counted: bool,
214}
215
216/// What one job of [`Merged::pages`] built.
217enum Built {
218    Stripe(ColumnStripe),
219    Block(EncodedBlock),
220}
221
222/// One column of a prepared stripe.
223#[derive(Debug)]
224enum Column {
225    /// Finished, because the column has no global dictionary.
226    Pages(ColumnStripe),
227    /// Coded against the stripe's own dictionary, waiting to be merged into the global one.
228    Coded(Local),
229}
230
231/// One column of a merged stripe.
232#[derive(Debug)]
233enum Merge {
234    Pages(ColumnStripe),
235    /// The local codes of every part, and the global code of every local one.
236    Codes {
237        parts: Vec<LocalPart>,
238        global: Vec<u32>,
239    },
240    /// A column that was prepared against a dictionary it no longer has, which is every column
241    /// prepared before the first stripe decided it should not have one. Encoded again, plainly,
242    /// from the values its stripe dictionary holds.
243    Plain(Local),
244}
245
246/// No value after this one has its hash.
247const END: u32 = u32::MAX;
248
249/// A dictionary of one column of one stripe.
250///
251/// The values are compared by their bytes rather than by a second hash, because they are all here
252/// to compare. The global dictionary has two hashes to go on because its values are mostly in the
253/// file by now. Both hashes are taken here, once a distinct value, so that merging it takes none.
254#[derive(Debug, Default)]
255struct Local {
256    /// The first value holding each hash.
257    first: HashMap<u64, u32, Spread>,
258    /// The next value holding the same hash as this one, or [`END`].
259    next: Vec<u32>,
260    hashes: Vec<u64>,
261    checks: Vec<u64>,
262    /// The values back to back, and where each one ends.
263    bytes: Vec<u8>,
264    ends: Vec<usize>,
265    /// How many rows that are not null hold each value, and how many are null.
266    counts: Vec<u64>,
267    nulls: u64,
268    parts: Vec<LocalPart>,
269    /// Whether the column is a blob rather than a varchar, for building its rows back.
270    blob: bool,
271}
272
273/// One part of one column coded against its stripe's dictionary.
274#[derive(Debug)]
275struct LocalPart {
276    codes: Vec<u32>,
277    /// What [`push_validity`] wrote for the part, which is the page's second field onwards.
278    validity: Vec<u8>,
279    range: Range,
280}
281
282impl Local {
283    /// One column of a stripe, coded in one go.
284    #[cfg(test)]
285    fn code_column(index: usize, held: &[PendingChunk]) -> Result<Self> {
286        let mut local = Self::default();
287        let mut mapped = None;
288        for pending in held {
289            local.code_part(pending.chunk.column(index)?, &mut mapped)?;
290        }
291        local.done();
292        Ok(local)
293    }
294
295    /// One more part of a column of a stripe, coded.
296    ///
297    /// A null row is coded as the empty string and counted as a null rather than against it, which
298    /// is what the writer has always done with one. The code is never read, since the page's
299    /// validity says the row is null, and giving it one keeps the page one code a row.
300    ///
301    /// The parts come to [`Local::code_part`] one at a time, in order, and [`Local::done`] ends the
302    /// column, so a stripe can be coded as its parts arrive rather than once they are all held.
303    fn code_part(
304        &mut self,
305        column: &Vector,
306        mapped: &mut Option<(Arc<Vector>, Vec<u32>)>,
307    ) -> Result<()> {
308        self.blob = column.logical_type() == &LogicalType::Blob;
309        if let Some(codes) = self.code_dictionary(column, mapped)? {
310            let mut validity = Vec::new();
311            push_validity(&mut validity, column);
312            self.parts.push(LocalPart { codes, validity, range: Range::of(column) });
313            return Ok(());
314        }
315        // flatten: the page is one code a row whatever form the rows came in.
316        let flat = column.flatten()?;
317        let mut codes = Vec::with_capacity(flat.len());
318        let mut last = None;
319        for row in 0..flat.len() {
320            // bytes_at rather than text_at: the rows were checked for UTF-8 when they came in,
321            // and checking every one again cost more than coding it.
322            let text = flat.bytes_at(row).unwrap_or(b"");
323            // A repeat of the row before is common enough on a sorted table to be worth a
324            // comparison before a hash, and the comparison fails on its first bytes when not.
325            let code = match last {
326                Some(code) if self.value(code) == text => code,
327                _ => self.code(text)?,
328            };
329            last = Some(code);
330            if flat.is_null_at(row) {
331                self.nulls += 1;
332            } else {
333                self.counts[code as usize] += 1;
334            }
335            codes.push(code);
336        }
337        let mut validity = Vec::new();
338        push_validity(&mut validity, &flat);
339        self.parts.push(LocalPart { codes, validity, range: Range::of(column) });
340        Ok(())
341    }
342
343    /// Ends a column once its last part is coded.
344    fn done(&mut self) {
345        // Only the coding needs to find a value by its bytes, and on a column of URLs the table
346        // that does it is as large as the codes.
347        self.first = HashMap::default();
348        self.next = Vec::new();
349    }
350
351    /// Codes a part that came in as codes into a dictionary of its own, which is how a Parquet page
352    /// written with a dictionary arrives, by coding each value of that dictionary once rather than
353    /// each row.
354    ///
355    /// Every row then costs a lookup in `mapped`, which holds the local code of each value of the
356    /// last dictionary seen by its position in it, and a value is coded the first time a row holds
357    /// it. So the codes come out in the order the rows first held each value, the same as coding the
358    /// rows one at a time, and the stripe is the same stripe either way. The parts of one Parquet
359    /// column chunk share their dictionary, so `mapped` carries over from one part to the next
360    /// while it is the same one, found by the pointer and kept alive by holding it.
361    ///
362    /// `None` for anything else, and for a dictionary with a null in it, since a null row there is
363    /// found through the value it points at and not through the part's own validity, which is the
364    /// one [`push_validity`] writes.
365    fn code_dictionary(
366        &mut self,
367        column: &Vector,
368        mapped: &mut Option<(Arc<Vector>, Vec<u32>)>,
369    ) -> Result<Option<Vec<u32>>> {
370        let Some((codes, values)) = column.shared_dictionary_parts() else { return Ok(None) };
371        if !matches!(values.validity(), Validity::AllValid) {
372            return Ok(None);
373        }
374        let Some(codes) = codes.get(..column.len()) else { return Ok(None) };
375        let fresh = !matches!(mapped, Some((held, _)) if Arc::ptr_eq(held, values));
376        if fresh {
377            *mapped = Some((Arc::clone(values), vec![END; values.len()]));
378        }
379        let Some((_, map)) = mapped.as_mut() else { return Ok(None) };
380        let every = matches!(column.validity(), Validity::AllValid);
381        let mut coded = Vec::with_capacity(codes.len());
382        for (row, &code) in codes.iter().enumerate() {
383            if !every && !column.validity().is_valid(row) {
384                // Coded as the empty string and counted as a null, as `code_part` does.
385                let code = self.code(b"")?;
386                self.nulls += 1;
387                coded.push(code);
388                continue;
389            }
390            let slot = map
391                .get_mut(code as usize)
392                .ok_or_else(|| invalid("a dictionary code is out of range"))?;
393            if *slot == END {
394                *slot = self.code(values.bytes_at(code as usize).unwrap_or(b""))?;
395            }
396            self.counts[*slot as usize] += 1;
397            coded.push(*slot);
398        }
399        Ok(Some(coded))
400    }
401
402    /// The column's parts as the rows they were coded from, for a column that lost its global
403    /// dictionary after this stripe was coded against one.
404    ///
405    /// A null row comes back as a null over the empty string, which is what it was coded as, and
406    /// each part gets back the same form of validity it had, since the page records which it was.
407    fn rows(&self) -> Result<Vec<Vector>> {
408        self.parts
409            .iter()
410            .map(|part| {
411                let len = part.codes.len();
412                let mut column = StringColumn::with_capacity(len);
413                for &code in &part.codes {
414                    column.push_bytes(self.value(code));
415                }
416                let validity = match part.validity.split_first() {
417                    Some((0, _)) => Validity::AllValid,
418                    Some((1, _)) => Validity::AllInvalid,
419                    Some((2, bits)) => {
420                        let mut mask = Bitmap::all_valid(len);
421                        for row in (0..len).filter(|row| bits[row / 8] & (1 << (row % 8)) == 0) {
422                            mask.set(row, false);
423                        }
424                        Validity::Mask(mask)
425                    }
426                    _ => return Err(Error::internal("a coded part has no validity")),
427                };
428                let ty = if self.blob { LogicalType::Blob } else { LogicalType::Varchar };
429                Ok(Vector::flat(ty, Data::Varlen(column))?.with_validity(validity))
430            })
431            .collect()
432    }
433
434    fn values(&self) -> usize {
435        self.ends.len()
436    }
437
438    fn value(&self, code: u32) -> &[u8] {
439        let code = code as usize;
440        let from = if code == 0 { 0 } else { self.ends[code - 1] };
441        &self.bytes[from..self.ends[code]]
442    }
443
444    fn code(&mut self, text: &[u8]) -> Result<u32> {
445        let hash = checksum(text);
446        let Some(&first) = self.first.get(&hash) else {
447            let code = self.push(text, hash)?;
448            self.first.insert(hash, code);
449            return Ok(code);
450        };
451        let mut at = first;
452        loop {
453            if self.value(at) == text {
454                return Ok(at);
455            }
456            match self.next[at as usize] {
457                END => break,
458                next => at = next,
459            }
460        }
461        let code = self.push(text, hash)?;
462        self.next[at as usize] = code;
463        Ok(code)
464    }
465
466    fn push(&mut self, text: &[u8], hash: u64) -> Result<u32> {
467        let code = u32::try_from(self.ends.len())
468            .ok()
469            .filter(|&code| code != END)
470            .ok_or_else(|| invalid("a stripe has too many values in one column"))?;
471        self.bytes.extend_from_slice(text);
472        self.ends.push(self.bytes.len());
473        self.next.push(END);
474        self.hashes.push(hash);
475        self.checks.push(seeded_checksum(text, DICTIONARY_CHECK_SEED));
476        self.counts.push(0);
477        Ok(code)
478    }
479
480    /// Puts every value into `dictionary` in the order this stripe first held it, and says what
481    /// each one's code is there.
482    fn merge_into(&self, dictionary: &mut GlobalDictionary) -> Result<Vec<u32>> {
483        let mut global = Vec::with_capacity(self.values());
484        for (code, (&hash, &check)) in self.hashes.iter().zip(&self.checks).enumerate() {
485            let text = self.value(code as u32);
486            let at = dictionary.code_hashed(text, hash, check)?;
487            let count = dictionary
488                .counts
489                .get_mut(at as usize)
490                .ok_or_else(|| invalid("global dictionary count code is out of range"))?;
491            *count = count.saturating_add(self.counts[code]);
492            global.push(at);
493        }
494        dictionary.nulls = dictionary.nulls.saturating_add(self.nulls);
495        Ok(global)
496    }
497}
498
499/// Whether a column's first stripe says it should not have a global dictionary.
500///
501/// Every varchar column starts with one, because the writer cannot know what is in a column before
502/// it has seen some of it. A global dictionary is the right shape for a column of a few dozen
503/// values repeated down the table: the pages become small integers, a filter against a literal is
504/// one search of the sorted order rather than a comparison a row, and a group by is on the codes.
505/// It is the wrong shape for a column whose values are nearly all different. There the codes are
506/// as wide as row numbers, nothing is saved on the pages, and the membership index of a stripe is
507/// a list of very nearly every code in the column. On TPC-H the orders table written on its own
508/// 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`
509/// from 1.810 G instructions to 1.213 G, which is what the rudb parquet reader takes over the same
510/// values.
511///
512/// So the first stripe of a column is the sample and the decision is made once on it. Once, rather
513/// than per stripe, because the codes of one column have to mean the same thing in every page of
514/// it, and a column that changed its mind halfway would need its earlier stripes rewritten. The
515/// first stripe is encoded again when the answer comes out against the dictionary, which is the one
516/// stripe that pays for the decision, along with any stripe that was prepared before it was made.
517///
518/// The threshold is deliberately near the top. [`DICTIONARY_DISTINCT_IN_TEN`] of the sample has to
519/// be values never seen before, which is a column with essentially no repeats. Everything with real
520/// repetition keeps its dictionary and keeps every property that hangs off it, and nothing is
521/// claimed here about where between the two the crossover really sits.
522fn drops_dictionary(rows: usize, distinct: usize) -> bool {
523    rows >= DICTIONARY_DECIDE_ROWS
524        && distinct.saturating_mul(10) > rows.saturating_mul(DICTIONARY_DISTINCT_IN_TEN)
525}
526
527/// Whether a column's dictionary should stop taking values after a stripe that added `new` of them
528/// in `rows` rows, with the dictionaries holding `total` bytes between them.
529///
530/// Section 5.5 of the encoding spec, which asks at every stripe what [`drops_dictionary`] asks at
531/// the first. A column that turns into a column of new values partway through, which is what a
532/// URL column of a log does once its first hours are past, stops growing its dictionary one stripe
533/// after it turns rather than at the end of the load. The stripes already coded keep their codes,
534/// so unlike the first stripe's decision this one costs nothing to make late.
535///
536/// The second reason is the cap. Once the dictionaries together hold more than it, the column that
537/// grew the most in its last stripe is the one that stops, because it is the one that would have
538/// taken the most of what is left.
539fn demotes(rows: usize, new: usize, total: u64, coding: &Coding, index: usize) -> bool {
540    drops_dictionary(rows, new)
541        || (total > coding.cap.load(Atomic::Relaxed) && coding.grew_most(index))
542}
543
544/// Runs `work` on every one of `jobs`, spread over `workers` threads, and hands back each job with
545/// what it came to, in no particular order.
546///
547/// The jobs are handed out through a queue rather than dealt in equal piles, because they are
548/// nothing like equal: `URL` on ClickBench is a string column of sixty one million distinct values
549/// and `IsMobile` is a byte. A pile that happened to hold the four large string columns would be
550/// the whole stripe and the other workers would be waiting on it. The caller hands the jobs over
551/// cheapest first and they are taken from the back, so the expensive ones go first, which is the
552/// classic answer to a last job that runs longer than everything before it.
553fn fan_out<T: Send>(
554    jobs: Vec<usize>,
555    workers: usize,
556    profile: Option<&LoadProfile>,
557    work: impl Fn(usize) -> Result<T> + Sync,
558) -> Result<Vec<(usize, T)>> {
559    if workers <= 1 || jobs.len() <= 1 {
560        let _span = profile.map(|profile| profile.span(Stage::Pages));
561        return jobs.into_iter().map(|index| Ok((index, work(index)?))).collect();
562    }
563    let workers = workers.min(jobs.len());
564    let queue = Mutex::new(jobs);
565    let pieces = std::thread::scope(|scope| {
566        (0..workers)
567            .map(|_| {
568                scope.spawn(|| {
569                    let _span = profile.map(|profile| profile.span(Stage::Pages));
570                    let mut mine = Vec::new();
571                    loop {
572                        let taken = queue
573                            .lock()
574                            .map_err(|_| Error::internal("a native encode worker panicked"))?
575                            .pop();
576                        let Some(index) = taken else { break };
577                        mine.push((index, work(index)?));
578                    }
579                    Ok(mine)
580                })
581            })
582            .collect::<Vec<_>>()
583            .into_iter()
584            .map(|handle| {
585                handle.join().map_err(|_| Error::internal("a native encode worker panicked"))?
586            })
587            .collect::<Result<Vec<Vec<_>>>>()
588    })?;
589    Ok(pieces.into_iter().flatten().collect())
590}
591
592impl Preparer {
593    /// Encodes a run of chunks as one stripe, as far as it can be without the writer.
594    ///
595    /// The run is what [`Writer::append_stripe`] takes, and the rules are the same: it is a stripe
596    /// of its own, and the orders have to come out in source order once the stripes are sorted. An
597    /// empty chunk is dropped.
598    ///
599    /// # Errors
600    ///
601    /// If the run is longer than [`STRIPE_PARTS`], a chunk's columns are not the table's, or one
602    /// cannot be encoded.
603    pub fn prepare(&self, parts: Vec<((u64, u64), Chunk)>) -> Result<Prepared> {
604        if parts.len() > STRIPE_PARTS {
605            return Err(invalid("a stripe was handed more parts than it holds"));
606        }
607        let held = parts
608            .into_iter()
609            .filter(|(_, chunk)| !chunk.is_empty())
610            .map(|(order, chunk)| PendingChunk { order, chunk })
611            .collect::<Vec<_>>();
612        for pending in &held {
613            self.fits(&pending.chunk)?;
614        }
615        self.prepare_held(held)
616    }
617
618    /// The check [`Writer::admit`] makes, here because the rows are gone by the merge.
619    fn fits(&self, chunk: &Chunk) -> Result<()> {
620        if chunk.width() != self.types.len() {
621            return Err(invalid("chunk width differs from table schema"));
622        }
623        for (index, ty) in self.types.iter().enumerate() {
624            if chunk.column(index)?.logical_type() != ty {
625                return Err(invalid("chunk type differs from table schema"));
626            }
627        }
628        Ok(())
629    }
630
631    pub(crate) fn prepare_held(&self, held: Vec<PendingChunk>) -> Result<Prepared> {
632        let mut building = self.start();
633        self.feed_held(&mut building, held)?;
634        self.finish(building)
635    }
636
637    /// Starts a stripe that will be handed over a few parts at a time.
638    ///
639    /// [`Preparer::prepare`] takes a stripe whole, which means a caller holds every part of it
640    /// decoded until the last one arrives, and every load instance holds one. Here each batch is
641    /// encoded as it comes and let go, so what an instance holds is one batch of rows and what it
642    /// has built from the batches before. The stripe comes out the same: the parts of a column go through
643    /// the same steps in the same order, only with the rows of later parts not yet in memory.
644    ///
645    /// Whether a column is coded against its global dictionary is read here, once, so every part of
646    /// the stripe is encoded the same way even if the column is demoted while it is being built.
647    /// A stripe that ends up coded against a dictionary its column no longer has is encoded again
648    /// plainly at the merge, which is what happens to a stripe prepared whole at the same moment.
649    #[must_use]
650    pub fn start(&self) -> Building {
651        let columns = (0..self.types.len())
652            .map(|index| {
653                let body = if self.coded[index].load(Atomic::Relaxed) {
654                    Body::Coded(Local::default(), None)
655                } else {
656                    Body::Pages(ColumnStripe::default(), Settling::default())
657                };
658                let mut gather = stats::Gather::new(&self.types[index], 0);
659                let ceiling = self.coded.ceilings[index].load(Atomic::Relaxed);
660                if let Some(gather) = gather.as_mut().filter(|_| ceiling < u64::MAX) {
661                    gather.cap_at(ceiling);
662                }
663                Mutex::new(Growing { body, gather })
664            })
665            .collect();
666        Building { parts: Vec::new(), columns }
667    }
668
669    /// Encodes the next few parts of a stripe that was started with [`Preparer::start`].
670    ///
671    /// The same rules as [`Preparer::prepare`]: the orders come in source order, an empty chunk is
672    /// dropped, and the stripe holds no more than [`STRIPE_PARTS`] in all.
673    ///
674    /// # Errors
675    ///
676    /// If the stripe would hold too many parts, a chunk's columns are not the table's, or one cannot
677    /// be encoded.
678    pub fn feed(&self, building: &mut Building, parts: Vec<((u64, u64), Chunk)>) -> Result<()> {
679        if building.parts.len().saturating_add(parts.len()) > STRIPE_PARTS {
680            return Err(invalid("a stripe was handed more parts than it holds"));
681        }
682        let held = parts
683            .into_iter()
684            .filter(|(_, chunk)| !chunk.is_empty())
685            .map(|(order, chunk)| PendingChunk { order, chunk })
686            .collect::<Vec<_>>();
687        for pending in &held {
688            self.fits(&pending.chunk)?;
689        }
690        self.feed_held(building, held)
691    }
692
693    fn feed_held(&self, building: &mut Building, held: Vec<PendingChunk>) -> Result<()> {
694        if held.is_empty() {
695            return Ok(());
696        }
697        let width = self.types.len();
698        // The stripe's key is its first part's order, and its statistics open with it.
699        let opening = building.parts.is_empty();
700        let key = held.first().map_or((0, 0), |pending| pending.order);
701        let share = Share::take(width, held.len());
702        let mut jobs = (0..width).collect::<Vec<_>>();
703        jobs.sort_by_key(|&index| weight(&self.types[index]));
704        let columns = &building.columns;
705        fan_out(jobs, share.0, self.profile.as_deref(), |index| {
706            let mut growing = columns[index]
707                .lock()
708                .map_err(|_| Error::internal("a native encode worker panicked"))?;
709            let Growing { body, gather } = &mut *growing;
710            if matches!(body, Body::Coded(..)) && !self.coded[index].load(Atomic::Relaxed) {
711                body.plain()?;
712            }
713            if let Some(gather) = gather.as_mut() {
714                // The statistics on the thread that is already walking the column, and in the same
715                // step, because the rows are in memory once and this is the moment they are.
716                if opening {
717                    gather.open_stripe(key);
718                }
719                for pending in &held {
720                    gather.part(pending.chunk.column(index)?);
721                }
722            }
723            match body {
724                Body::Coded(local, mapped) => {
725                    for pending in &held {
726                        local.code_part(pending.chunk.column(index)?, mapped)?;
727                    }
728                }
729                Body::Pages(stripe, settling) => {
730                    for pending in &held {
731                        Writer::encode_page(stripe, settling, pending.chunk.column(index)?)?;
732                    }
733                }
734            }
735            Ok(())
736        })?;
737        drop(share);
738        building.parts.extend(held.iter().map(Part::of));
739        Ok(())
740    }
741
742    /// Ends a stripe that was started with [`Preparer::start`], ready for the merge.
743    ///
744    /// # Errors
745    ///
746    /// If a column's worker panicked.
747    pub fn finish(&self, building: Building) -> Result<Prepared> {
748        let empty = building.parts.is_empty();
749        let (columns, gathers) = building
750            .columns
751            .into_iter()
752            .enumerate()
753            .map(|(index, growing)| {
754                let Growing { body, gather } = growing
755                    .into_inner()
756                    .map_err(|_| Error::internal("a native encode worker panicked"))?;
757                let column = match body {
758                    Body::Coded(mut local, _) => {
759                        local.done();
760                        Column::Coded(local)
761                    }
762                    Body::Pages(stripe, _) => Column::Pages(stripe),
763                };
764                // A stripe of no parts has no statistics, rather than an empty stripe of them.
765                let gather = gather.filter(|_| !empty).map(|mut gather| {
766                    gather.close_stripe();
767                    if let Some(ceiling) = gather.ceiling() {
768                        self.coded.ceilings[index].fetch_min(ceiling, Atomic::Relaxed);
769                    }
770                    gather
771                });
772                Ok((column, gather))
773            })
774            .collect::<Result<Vec<_>>>()?
775            .into_iter()
776            .unzip();
777        Ok(Prepared {
778            parts: building.parts,
779            types: self.types.clone(),
780            columns,
781            gathers,
782            profile: self.profile.clone(),
783        })
784    }
785}
786
787/// A stripe that [`Preparer::start`] began and [`Preparer::feed`] is adding parts to.
788#[derive(Debug)]
789pub struct Building {
790    parts: Vec<Part>,
791    /// Behind a lock each so a column can be handed to whichever worker takes it. Only one does
792    /// at a time, so the locks are never waited on.
793    columns: Vec<Mutex<Growing>>,
794}
795
796impl Building {
797    /// How many parts the stripe holds so far.
798    #[must_use]
799    pub fn parts(&self) -> usize {
800        self.parts.len()
801    }
802}
803
804/// One column of a stripe that is being built.
805#[derive(Debug)]
806struct Growing {
807    body: Body,
808    gather: Option<stats::Gather>,
809}
810
811/// What one column of a stripe being built has come to so far.
812#[derive(Debug)]
813enum Body {
814    /// Pages, with what the parts so far have settled on.
815    Pages(ColumnStripe, Settling),
816    /// Codes against the stripe's own dictionary, with the Parquet dictionary the last part came
817    /// in, as [`Local::code_dictionary`] keeps it.
818    Coded(Local, Option<(Arc<Vector>, Vec<u32>)>),
819}
820
821impl Body {
822    /// Turns a column coded against its dictionary into pages, for a column that lost its
823    /// dictionary while the stripe was being built.
824    ///
825    /// The merge would encode the whole stripe again plainly once it saw the column had none. Doing
826    /// it here, on the parts coded so far, means the parts still to come are encoded once. A load
827    /// starts every stripe it has in flight before the first one reaches the merge and decides.
828    fn plain(&mut self) -> Result<()> {
829        let Self::Coded(local, _) = self else { return Ok(()) };
830        let mut stripe = ColumnStripe::default();
831        let mut settling = Settling::default();
832        for rows in local.rows()? {
833            Writer::encode_page(&mut stripe, &mut settling, &rows)?;
834        }
835        *self = Self::Pages(stripe, settling);
836        Ok(())
837    }
838}
839
840/// Where one column's dictionary and statistics are while a stripe is merged into them.
841enum Slot<'a> {
842    /// In the writer, which the caller holds.
843    Owned(&'a mut Option<GlobalDictionary>, &'a mut Option<stats::Gather>),
844    /// Lent to a [`Merger`], behind the column's own lock.
845    Lent(&'a Mutex<LentColumn>, &'a Lent),
846}
847
848/// One column of a stripe on its way through [`merge_columns`].
849struct Step<'a> {
850    index: usize,
851    column: Column,
852    slot: Slot<'a>,
853    /// The stripe's statistics for the column, when the column keeps them.
854    gather: Option<stats::Gather>,
855}
856
857impl Step<'_> {
858    /// Roughly what the merge costs: a hash a distinct value when there is a global dictionary to
859    /// merge into, and next to nothing otherwise. Read off `coded` rather than the dictionary, so a
860    /// lent column does not have to be locked to be sorted.
861    fn cost(&self, coded: &Coding) -> usize {
862        match &self.column {
863            Column::Coded(local) if coded[self.index].load(Atomic::Relaxed) => {
864                local.values().saturating_add(1)
865            }
866            _ => 0,
867        }
868    }
869
870    /// Merges the column, settles its dictionary's shape and hands out the blocks it filled.
871    fn run(
872        self,
873        rows: usize,
874        coded: &Coding,
875        profile: Option<&LoadProfile>,
876    ) -> Result<(usize, Merge, Vec<Unencoded>)> {
877        let Self { index, column, slot, gather } = self;
878        match slot {
879            Slot::Owned(dictionary, mine) => {
880                merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
881            }
882            Slot::Lent(held, lent) => {
883                let mut held = held.lock().map_err(|_| Error::internal("a merge panicked"))?;
884                // Checked with the column locked, so a merge either finishes before the writer
885                // takes this column back or is refused.
886                if lent.reclaimed.load(Atomic::Acquire) {
887                    return Err(Error::internal("a stripe was merged after its table was closed"));
888                }
889                let LentColumn { dictionary, gather: mine } = &mut *held;
890                merge_column(index, column, gather, dictionary, mine, rows, coded, profile)
891            }
892        }
893    }
894}
895
896/// One column of [`merge_columns`].
897#[expect(clippy::too_many_arguments, reason = "one column's share of the stripe's merge state")]
898fn merge_column(
899    index: usize,
900    column: Column,
901    stripe: Option<stats::Gather>,
902    dictionary: &mut Option<GlobalDictionary>,
903    gather: &mut Option<stats::Gather>,
904    rows: usize,
905    coded: &Coding,
906    profile: Option<&LoadProfile>,
907) -> Result<(usize, Merge, Vec<Unencoded>)> {
908    if let (Some(mine), Some(stripe)) = (gather.as_mut(), stripe) {
909        mine.absorb(stripe);
910    }
911    let mut new = None;
912    let merge = match (column, dictionary.as_mut()) {
913        (Column::Pages(stripe), None) => Merge::Pages(stripe),
914        (Column::Pages(stripe), Some(global)) if global.demoted => Merge::Pages(stripe),
915        (Column::Pages(_), Some(_)) => {
916            return Err(Error::internal(
917                "a column with a global dictionary was prepared without one",
918            ));
919        }
920        (Column::Coded(local), None) => Merge::Plain(local),
921        // Prepared before the column was demoted and merged after.
922        (Column::Coded(local), Some(global)) if global.demoted => Merge::Plain(local),
923        (Column::Coded(local), Some(global)) => {
924            // Empty means nothing has been merged into it yet, so this is the column's first
925            // stripe and the only one the decision is allowed to be made on.
926            if global.values() == 0 && drops_dictionary(rows, local.values()) {
927                if let Some(profile) = profile {
928                    profile.release(global.charged);
929                }
930                coded.recount(global.charged, 0);
931                *dictionary = None;
932                coded[index].store(false, Atomic::Relaxed);
933                Merge::Plain(local)
934            } else {
935                let before = global.values();
936                let codes = local.merge_into(global)?;
937                new = Some(global.values() - before);
938                Merge::Codes { parts: local.parts, global: codes }
939            }
940        }
941    };
942    // Asked before the blocks go out, so that a demotion's sealed part block goes out with them.
943    if let (Some(new), Some(global)) = (new, dictionary.as_mut()) {
944        let now = global.held_bytes();
945        coded.growth[index].store(now.saturating_sub(global.charged), Atomic::Relaxed);
946        let total =
947            coded.held.load(Atomic::Relaxed).saturating_add(now).saturating_sub(global.charged);
948        if demotes(rows, new, total, coded, index) {
949            global.demote();
950            coded[index].store(false, Atomic::Relaxed);
951            coded.growth[index].store(0, Atomic::Relaxed);
952        }
953    }
954    // Settled here rather than when the stripe is written, so that the blocks this merge filled go
955    // out with it already knowing their shape. A column still too small to settle one keeps its
956    // blocks until it can, which is at most `PAYLOAD_SAMPLE_BLOCKS` of them, because encoding them
957    // now would be encoding them without having looked at the column.
958    let blocks = match dictionary {
959        Some(dictionary) => {
960            dictionary.settle()?;
961            let blocks = dictionary.hand_out(index);
962            let (before, now) = dictionary.recharge(profile);
963            coded.recount(before, now);
964            blocks
965        }
966        None => Vec::new(),
967    };
968    Ok((index, merge, blocks))
969}
970
971/// Merges every column of a stripe into the dictionaries and statistics in `slots`.
972///
973/// Every column is merged on its own, because nothing one column's merge reads or writes belongs to
974/// another: its statistics, its global dictionary and its flag in `coded`. So the columns are
975/// spread over threads, and a stripe takes as long as its slowest column rather than all of them.
976/// The answer is the same in any order, because a column's merge only depends on the stripes
977/// merged into that column before it.
978fn merge_columns(prepared: Prepared, slots: Vec<Slot<'_>>, coded: &Coding) -> Result<Merged> {
979    let Prepared { parts, columns, gathers, profile, .. } = prepared;
980    let timing = profile.as_deref().map(|profile| profile.span(Stage::Dictionary));
981    let rows: usize = parts.iter().map(|part| part.rows).sum();
982    let width = columns.len();
983    if slots.len() != width || gathers.len() != width {
984        return Err(Error::internal("a stripe was merged into a table of another width"));
985    }
986    let mut steps = columns
987        .into_iter()
988        .zip(gathers)
989        .zip(slots)
990        .enumerate()
991        .map(|(index, ((column, gather), slot))| Step { index, column, slot, gather })
992        .collect::<Vec<_>>();
993    // Taken from the back, so the biggest merges start first and the last one to finish is
994    // small, the same reason `fan_out` hands its jobs over cheapest first.
995    steps.sort_by_key(|step| step.cost(coded));
996    let workers = std::thread::available_parallelism()
997        .map_or(1, usize::from)
998        .min(MAX_ENCODE_WORKERS)
999        .min(steps.iter().filter(|step| step.cost(coded) > 0).count())
1000        .max(1);
1001    let done = if workers <= 1 {
1002        steps
1003            .into_iter()
1004            .map(|step| step.run(rows, coded, profile.as_deref()))
1005            .collect::<Result<Vec<_>>>()?
1006    } else {
1007        let queue = Mutex::new(steps);
1008        let pieces = std::thread::scope(|scope| {
1009            (0..workers)
1010                .map(|_| {
1011                    scope.spawn(|| {
1012                        let mut mine = Vec::new();
1013                        loop {
1014                            let taken = queue
1015                                .lock()
1016                                .map_err(|_| Error::internal("a merge worker panicked"))?
1017                                .pop();
1018                            let Some(step) = taken else { break };
1019                            mine.push(step.run(rows, coded, profile.as_deref())?);
1020                        }
1021                        Ok(mine)
1022                    })
1023                })
1024                .collect::<Vec<_>>()
1025                .into_iter()
1026                .map(|handle| {
1027                    handle.join().map_err(|_| Error::internal("a merge worker panicked"))?
1028                })
1029                .collect::<Result<Vec<Vec<_>>>>()
1030        })?;
1031        pieces.into_iter().flatten().collect()
1032    };
1033    let mut slots: Vec<Option<(Merge, Vec<Unencoded>)>> = (0..width).map(|_| None).collect();
1034    for (index, merge, blocks) in done {
1035        slots[index] = Some((merge, blocks));
1036    }
1037    let mut merged = Vec::with_capacity(width);
1038    let mut blocks = Vec::new();
1039    for slot in slots {
1040        let (merge, handed) = slot.ok_or_else(|| Error::internal("a column was never merged"))?;
1041        merged.push(merge);
1042        blocks.extend(handed);
1043    }
1044    drop(timing);
1045    Ok(Merged { parts, columns: merged, blocks, profile, counted: false })
1046}
1047
1048/// The dictionaries and statistics of a table while a [`Merger`] has them, one lock a column.
1049#[derive(Debug)]
1050pub(crate) struct Lent {
1051    columns: Box<[Mutex<LentColumn>]>,
1052    /// Set when the writer takes them back, after which a merge is refused.
1053    reclaimed: AtomicBool,
1054}
1055
1056/// One column of [`Lent`].
1057#[derive(Debug)]
1058pub(crate) struct LentColumn {
1059    pub(crate) dictionary: Option<GlobalDictionary>,
1060    gather: Option<stats::Gather>,
1061}
1062
1063impl Lent {
1064    pub(crate) fn columns(&self) -> &[Mutex<LentColumn>] {
1065        &self.columns
1066    }
1067
1068    /// Puts encoded blocks back into their dictionaries, each under its own column's lock.
1069    fn take_back(&self, blocks: Vec<(usize, usize, EncodedBlock)>) -> Result<()> {
1070        for (column, at, block) in blocks {
1071            self.columns
1072                .get(column)
1073                .ok_or_else(|| Error::internal("a dictionary block came back to no column"))?
1074                .lock()
1075                .map_err(|_| Error::internal("a merge panicked"))?
1076                .dictionary
1077                .as_mut()
1078                .ok_or_else(|| Error::internal("a dictionary block came back to no dictionary"))?
1079                .take_back(at, block)?;
1080        }
1081        Ok(())
1082    }
1083
1084    /// Everything lent, handed back to the writer.
1085    #[allow(clippy::type_complexity)]
1086    pub(crate) fn reclaim(
1087        &self,
1088    ) -> Result<(Vec<Option<GlobalDictionary>>, Vec<Option<stats::Gather>>)> {
1089        self.reclaimed.store(true, Atomic::Release);
1090        let mut dictionaries = Vec::with_capacity(self.columns.len());
1091        let mut gathers = Vec::with_capacity(self.columns.len());
1092        for column in &self.columns {
1093            let mut held = column.lock().map_err(|_| Error::internal("a merge panicked"))?;
1094            dictionaries.push(held.dictionary.take());
1095            gathers.push(held.gather.take());
1096        }
1097        Ok((dictionaries, gathers))
1098    }
1099}
1100
1101/// Merges prepared stripes into a writer's dictionaries and statistics without the writer.
1102///
1103/// Handed out by [`Writer::merger`]. With it, a load that shares one writer between many threads
1104/// holds the writer's lock only to write, and two stripes merge at once as long as they are on
1105/// different columns. A stripe merged here is written with [`Writer::write`] as usual, and that is
1106/// where its rows are counted in.
1107#[derive(Debug, Clone)]
1108pub struct Merger {
1109    lent: Arc<Lent>,
1110    types: Vec<LogicalType>,
1111    coded: Arc<Coding>,
1112}
1113
1114impl Merger {
1115    /// [`Writer::merge`], one column lock at a time instead of the writer.
1116    ///
1117    /// # Errors
1118    ///
1119    /// If the stripe was prepared for a table of other columns, or the table was closed.
1120    pub fn merge(&self, prepared: Prepared) -> Result<Merged> {
1121        if prepared.types != self.types {
1122            return Err(invalid("a stripe was prepared for a table of other columns"));
1123        }
1124        let slots = self.lent.columns.iter().map(|column| Slot::Lent(column, &self.lent)).collect();
1125        merge_columns(prepared, slots, &self.coded)
1126    }
1127
1128    /// Puts a stripe's encoded dictionary blocks back, so that [`Writer::write`] does not wait on a
1129    /// column's lock while it holds its own.
1130    ///
1131    /// # Errors
1132    ///
1133    /// If a block comes back to a column without a dictionary, or comes back twice.
1134    pub fn give_back(&self, paged: &mut Paged) -> Result<()> {
1135        self.lent.take_back(std::mem::take(&mut paged.blocks))
1136    }
1137}
1138
1139impl Merged {
1140    /// Builds the pages the merge left to build, which is every column coded against a global
1141    /// dictionary and every column that lost one after the stripe was prepared, and encodes the
1142    /// dictionary blocks the merge filled.
1143    ///
1144    /// # Errors
1145    ///
1146    /// If a column or a block cannot be encoded or a page comes out larger than a page may be.
1147    pub fn pages(self) -> Result<Paged> {
1148        let Self { parts, columns, blocks, profile, counted } = self;
1149        let width = columns.len();
1150        // The blocks go first so that they are taken last. One block is a thousand values, which is
1151        // less than any column of a stripe, and small jobs at the end are what keeps the last
1152        // worker from finishing long after the others.
1153        let mut jobs = (width..width + blocks.len())
1154            .chain((0..width).filter(|&index| !matches!(columns[index], Merge::Pages(_))))
1155            .collect::<Vec<_>>();
1156        // A column encoded again from its rows costs more than one whose codes only need building.
1157        jobs.sort_by_key(|&index| index < width && matches!(columns[index], Merge::Plain(_)));
1158        let share = Share::take(jobs.len(), parts.len());
1159        let built = fan_out(jobs, share.0, profile.as_deref(), |index| {
1160            let Some(column) = columns.get(index) else {
1161                return Ok(Built::Block(blocks[index - width].encode()?));
1162            };
1163            Ok(Built::Stripe(match column {
1164                Merge::Codes { parts, global } => code_pages(parts, global)?,
1165                Merge::Plain(local) => {
1166                    Writer::encode_pages(&local.rows()?.iter().collect::<Vec<_>>())?
1167                }
1168                Merge::Pages(_) => {
1169                    return Err(Error::internal("a finished column was queued to be built"));
1170                }
1171            }))
1172        })?;
1173        drop(share);
1174        let mut slots: Vec<Option<ColumnStripe>> = (0..width).map(|_| None).collect();
1175        let mut encoded = Vec::with_capacity(blocks.len());
1176        for (index, one) in built {
1177            match one {
1178                Built::Stripe(stripe) => slots[index] = Some(stripe),
1179                Built::Block(block) => {
1180                    let (column, at) = blocks[index - width].place();
1181                    encoded.push((column, at, block));
1182                }
1183            }
1184        }
1185        let columns = columns
1186            .into_iter()
1187            .zip(slots)
1188            .map(|(column, slot)| match (column, slot) {
1189                (Merge::Pages(stripe), _) | (_, Some(stripe)) => Ok(stripe),
1190                _ => Err(Error::internal("a column was never encoded")),
1191            })
1192            .collect::<Result<Vec<_>>>()?;
1193        Ok(Paged { parts, columns, blocks: encoded, counted })
1194    }
1195}
1196
1197/// One column's parts as pages of global codes.
1198fn code_pages(parts: &[LocalPart], global: &[u32]) -> Result<ColumnStripe> {
1199    let mut stripe = ColumnStripe {
1200        pages: Vec::with_capacity(parts.len()),
1201        codes: Vec::with_capacity(parts.len()),
1202        sieves: Vec::with_capacity(parts.len()),
1203        ranges: Vec::with_capacity(parts.len()),
1204    };
1205    for part in parts {
1206        let codes = part
1207            .codes
1208            .iter()
1209            .map(|&code| global.get(code as usize).copied())
1210            .collect::<Option<Vec<_>>>()
1211            .ok_or_else(|| Error::internal("a stripe's code has no global code"))?;
1212        let bytes = coded_page(&codes, &part.validity)?;
1213        if bytes.len() > MAX_PAGE {
1214            return Err(invalid("column page exceeds the configured bound"));
1215        }
1216        stripe.pages.push(bytes);
1217        stripe.codes.push(Some(unique_codes(&codes)));
1218        // None, because the codes already give the stripe an exact membership index, and an
1219        // approximate one beside it would cost a hash of every string to answer a question that
1220        // is already answered.
1221        stripe.sieves.push(None);
1222        stripe.ranges.push(part.range.clone());
1223    }
1224    Ok(stripe)
1225}
1226
1227impl Writer {
1228    /// Something that encodes stripes for this writer without holding it. See [`Preparer::prepare`].
1229    ///
1230    /// It carries the profile the writer has when it is asked for, so a writer that is going to be
1231    /// given one with [`Writer::with_profile`] should be given it first.
1232    #[must_use]
1233    pub fn preparer(&self) -> Preparer {
1234        Preparer {
1235            types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1236            coded: Arc::clone(&self.coded),
1237            profile: self.profile.clone(),
1238        }
1239    }
1240
1241    /// Takes a prepared stripe into the table's dictionaries and statistics and counts its rows in.
1242    ///
1243    /// This is the step that has to see the stripes one at a time, and it is a hash a distinct
1244    /// value of each varchar column rather than two a row. Whatever [`Writer::append_at`] left
1245    /// behind is written first as its own stripe, the same rule [`Writer::append_stripe`] has.
1246    ///
1247    /// # Errors
1248    ///
1249    /// If the stripe was prepared for a table of other columns, or the buffered stripe cannot be
1250    /// written.
1251    pub fn merge(&mut self, prepared: Prepared) -> Result<Merged> {
1252        self.flush_pending()?;
1253        if prepared.columns.len() != self.table.fields.len()
1254            || prepared.types.iter().ne(self.table.fields.iter().map(|field| &field.ty))
1255        {
1256            return Err(invalid("a stripe was prepared for a table of other columns"));
1257        }
1258        self.table.rows = prepared
1259            .parts
1260            .iter()
1261            .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1262            .ok_or_else(|| invalid("row count overflow"))?;
1263        self.merge_held(prepared)
1264    }
1265
1266    /// [`Writer::merge`] for a stripe whose rows are already counted in.
1267    ///
1268    /// Every column is merged on its own, because nothing one column's merge reads or writes
1269    /// belongs to another: its statistics, its global dictionary and its flag in `coded`. So the
1270    /// columns are spread over threads, and the lock is held for the slowest column rather than for
1271    /// all of them. On ClickBench `hits` the lock was busy 98% of a load and the merge was four
1272    /// fifths of that, while two thirds of the machine waited for it. The answer is the same in any
1273    /// order, because a column's merge only depends on the stripes merged into it before.
1274    pub(crate) fn merge_held(&mut self, prepared: Prepared) -> Result<Merged> {
1275        let slots = match &self.lent {
1276            Some(lent) => lent.columns.iter().map(|column| Slot::Lent(column, lent)).collect(),
1277            None => self
1278                .dictionaries
1279                .iter_mut()
1280                .zip(self.gathers.iter_mut())
1281                .map(|(dictionary, gather)| Slot::Owned(dictionary, gather))
1282                .collect::<Vec<_>>(),
1283        };
1284        let mut merged = merge_columns(prepared, slots, &self.coded)?;
1285        merged.counted = true;
1286        Ok(merged)
1287    }
1288
1289    /// Hands the dictionaries and the statistics to a [`Merger`], so that stripes can be merged
1290    /// without this writer's lock.
1291    ///
1292    /// Whatever [`Writer::append_at`] left behind is written first, the same rule
1293    /// [`Writer::merge`] has. The writer takes them back when the table is closed.
1294    ///
1295    /// # Errors
1296    ///
1297    /// If the buffered stripe cannot be written.
1298    pub fn merger(&mut self) -> Result<Merger> {
1299        self.flush_pending()?;
1300        let lent = match &self.lent {
1301            Some(lent) => Arc::clone(lent),
1302            None => {
1303                let lent = Arc::new(Lent {
1304                    columns: std::mem::take(&mut self.dictionaries)
1305                        .into_iter()
1306                        .zip(std::mem::take(&mut self.gathers))
1307                        .map(|(dictionary, gather)| Mutex::new(LentColumn { dictionary, gather }))
1308                        .collect(),
1309                    reclaimed: AtomicBool::new(false),
1310                });
1311                self.lent = Some(Arc::clone(&lent));
1312                lent
1313            }
1314        };
1315        Ok(Merger {
1316            lent,
1317            types: self.table.fields.iter().map(|field| field.ty.clone()).collect(),
1318            coded: Arc::clone(&self.coded),
1319        })
1320    }
1321
1322    /// Writes a stripe whose pages are built.
1323    ///
1324    /// # Errors
1325    ///
1326    /// If the stripe was built for a table of another width or cannot be written.
1327    pub fn write(&mut self, paged: Paged) -> Result<()> {
1328        self.write_paged(paged)
1329    }
1330
1331    pub(crate) fn write_paged(&mut self, paged: Paged) -> Result<()> {
1332        let Paged { parts, columns, blocks, counted } = paged;
1333        if !counted {
1334            self.table.rows = parts
1335                .iter()
1336                .try_fold(self.table.rows, |rows, part| rows.checked_add(part.rows))
1337                .ok_or_else(|| invalid("row count overflow"))?;
1338        }
1339        if let Some(lent) = &self.lent {
1340            lent.take_back(blocks)?;
1341        } else {
1342            for (column, at, block) in blocks {
1343                self.dictionaries
1344                    .get_mut(column)
1345                    .and_then(Option::as_mut)
1346                    .ok_or_else(|| {
1347                        Error::internal("a dictionary block came back to no dictionary")
1348                    })?
1349                    .take_back(at, block)?;
1350            }
1351        }
1352        if parts.is_empty() {
1353            return self.place_blocks();
1354        }
1355        self.write_stripe(&parts, columns)
1356    }
1357
1358    /// All four steps one after the other, for a caller with nobody to share the writer with.
1359    ///
1360    /// # Errors
1361    ///
1362    /// The same as [`Writer::merge`], [`Merged::pages`] and [`Writer::write`].
1363    pub fn append_prepared(&mut self, prepared: Prepared) -> Result<()> {
1364        let merged = self.merge(prepared)?;
1365        let paged = merged.pages()?;
1366        self.write(paged)
1367    }
1368}
1369
1370#[cfg(test)]
1371mod tests {
1372    use std::fs;
1373    use std::path::PathBuf;
1374    use std::time::{SystemTime, UNIX_EPOCH};
1375
1376    use rudb_common::{Field, Value};
1377    use rudb_vector::Vector;
1378
1379    use super::*;
1380    use crate::Reader;
1381
1382    const PART: usize = 1_000;
1383
1384    /// A stripe that came in as codes into Parquet style dictionaries codes to the same stripe as
1385    /// the same rows would flat: the same values in the same order, the same codes, counts, nulls
1386    /// and validity. Two of the parts share a dictionary and one has one of its own, and one of the
1387    /// shared ones has nulls pointing at a value that is not the empty string. None of them is
1388    /// flattened on the way.
1389    #[test]
1390    fn dictionary_parts_code_as_their_rows_would() {
1391        let texts = |values: &[&str]| {
1392            Arc::new(
1393                Vector::from_values(
1394                    LogicalType::Varchar,
1395                    &values
1396                        .iter()
1397                        .map(|text| Value::Varchar((*text).to_string()))
1398                        .collect::<Vec<_>>(),
1399                )
1400                .expect("a dictionary"),
1401            )
1402        };
1403        let shared = texts(&["b", "a", "", "c", "unused"]);
1404        let other = texts(&["c", "d", "a"]);
1405        let mut nulls = Bitmap::all_valid(6);
1406        nulls.set(1, false);
1407        nulls.set(4, false);
1408        let parts = [
1409            Vector::dictionary_over(vec![3, 3, 1, 0, 2, 1], Arc::clone(&shared)).expect("codes"),
1410            Vector::dictionary_over(vec![0, 3, 1, 1, 3, 2], Arc::clone(&shared))
1411                .expect("codes")
1412                .with_validity(Validity::Mask(nulls)),
1413            Vector::dictionary_over(vec![1, 2, 0, 1], other).expect("codes"),
1414        ];
1415        let held = |flat: bool| {
1416            parts
1417                .iter()
1418                .enumerate()
1419                .map(|(at, part)| PendingChunk {
1420                    order: (at as u64, 0),
1421                    chunk: Chunk::new(vec![if flat {
1422                        part.flatten().expect("flat")
1423                    } else {
1424                        part.clone()
1425                    }])
1426                    .expect("a chunk"),
1427                })
1428                .collect::<Vec<_>>()
1429        };
1430        let parquet = held(false);
1431        let before = rudb_common::slow::here();
1432        let coded = Local::code_column(0, &parquet).expect("coded");
1433        assert_eq!(
1434            rudb_common::slow::here().since(before).get(rudb_common::slow::Cause::Flatten),
1435            0,
1436            "a part that came in as codes was flattened",
1437        );
1438        let flat = Local::code_column(0, &held(true)).expect("coded");
1439        assert_eq!(coded.values(), flat.values());
1440        for code in 0..flat.values() as u32 {
1441            assert_eq!(coded.value(code), flat.value(code), "value {code}");
1442        }
1443        assert_eq!(coded.counts, flat.counts);
1444        assert_eq!(coded.nulls, flat.nulls);
1445        assert_eq!(coded.nulls, 2);
1446        assert_eq!(coded.hashes, flat.hashes);
1447        assert_eq!(coded.checks, flat.checks);
1448        assert_eq!(coded.parts.len(), flat.parts.len());
1449        for (coded, flat) in coded.parts.iter().zip(&flat.parts) {
1450            assert_eq!(coded.codes, flat.codes);
1451            assert_eq!(coded.validity, flat.validity);
1452        }
1453    }
1454
1455    fn path(label: &str) -> PathBuf {
1456        let stamp = SystemTime::now().duration_since(UNIX_EPOCH).expect("time advances").as_nanos();
1457        std::env::temp_dir()
1458            .join(format!("rudb-prepare-{label}-{}-{stamp}.rdb", std::process::id()))
1459    }
1460
1461    fn fields() -> Vec<Field> {
1462        vec![
1463            Field::required("id", LogicalType::BigInt),
1464            Field::new("city", LogicalType::Varchar),
1465            Field::new("note", LogicalType::Varchar),
1466        ]
1467    }
1468
1469    /// The value every row holds, so a test can check a row it reads back without keeping the rows.
1470    ///
1471    /// `city` repeats a handful of values and has a null every so often, which keeps its dictionary.
1472    /// `note` is different on every row but its nulls, which loses it on the first stripe.
1473    fn row(id: usize) -> [Value; 3] {
1474        let city = if id % 11 == 0 {
1475            Value::Null
1476        } else {
1477            Value::Varchar(format!("city {}", (id / 7) % 13))
1478        };
1479        let note = if id % 17 == 0 { Value::Null } else { Value::Varchar(format!("note {id}")) };
1480        [Value::BigInt(id as i64), city, note]
1481    }
1482
1483    /// A run of `parts` chunks starting at part `first`, as a caller hands them to the writer.
1484    fn stripe(first: usize, parts: usize) -> Vec<((u64, u64), Chunk)> {
1485        (first..first + parts)
1486            .map(|part| {
1487                let rows = (part * PART..(part + 1) * PART).map(row).collect::<Vec<_>>();
1488                let column = |at: usize| {
1489                    let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1490                    Vector::from_values(fields()[at].ty.clone(), &values).expect("a column")
1491                };
1492                let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1493                ((part as u64, 0), chunk)
1494            })
1495            .collect()
1496    }
1497
1498    /// The runs the tests hand over, out of source order so that the stripes are sorted at commit.
1499    fn runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1500        vec![stripe(5, 5), stripe(0, 5), stripe(10, 3)]
1501    }
1502
1503    fn check(path: &PathBuf) {
1504        let reader = Reader::open(path).expect("reopen");
1505        assert_eq!(reader.parts(), 13);
1506        for part in 0..13 {
1507            let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1508            for at in [0, 17, PART - 1] {
1509                let want = row(part * PART + at);
1510                for (column, value) in want.iter().enumerate() {
1511                    assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1512                }
1513            }
1514        }
1515    }
1516
1517    /// Every stripe prepared before any of them is merged writes the file that handing the same
1518    /// runs to the writer one at a time writes, byte for byte.
1519    ///
1520    /// That is the claim the whole split rests on. The second and third stripes here are coded
1521    /// against a dictionary for `note`, which the first stripe to be merged then decides the column
1522    /// should not have, so they are encoded again without it. `city` keeps its dictionary and the
1523    /// later stripes' values go into it in the order the merges happen.
1524    #[test]
1525    fn stripes_prepared_before_any_is_merged_write_the_same_bytes_as_one_at_a_time() {
1526        let alone = path("alone");
1527        let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1528        for run in runs() {
1529            writer.append_stripe(run).expect("a stripe");
1530        }
1531        writer.finish().expect("commit");
1532
1533        let split = path("split");
1534        let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1535        let preparer = writer.preparer();
1536        let prepared = runs()
1537            .into_iter()
1538            .map(|run| preparer.prepare(run).expect("prepared"))
1539            .collect::<Vec<_>>();
1540        for one in prepared {
1541            writer.append_prepared(one).expect("a stripe");
1542        }
1543        assert!(!preparer.coded[2].load(Atomic::Relaxed), "note lost its dictionary");
1544        assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1545        writer.finish().expect("commit");
1546
1547        assert_eq!(fs::read(&alone).expect("read"), fs::read(&split).expect("read"));
1548        check(&split);
1549        fs::remove_file(alone).expect("remove");
1550        fs::remove_file(split).expect("remove");
1551    }
1552
1553    /// Stripes counted under the ceiling an earlier stripe reported come to the same distinct count
1554    /// as stripes counted from nothing, whatever order they are merged in.
1555    ///
1556    /// The file cannot say this on its own, because a table this small may not be given the budget
1557    /// for its sketches, so the writer's statistics are compared before it closes.
1558    #[test]
1559    fn stripes_counted_under_an_earlier_ceiling_count_what_they_would_have() {
1560        let estimates = |writer: &Writer| {
1561            writer
1562                .gathers
1563                .iter()
1564                .map(|gather| gather.as_ref().and_then(stats::Gather::distinct))
1565                .collect::<Vec<_>>()
1566        };
1567        // What one gather makes of every row, with no stripe and so no ceiling anywhere.
1568        let want = fields()
1569            .iter()
1570            .enumerate()
1571            .map(|(column, field)| {
1572                let mut gather = stats::Gather::new(&field.ty, 0)?;
1573                for (_, chunk) in runs().into_iter().flatten() {
1574                    gather.part(chunk.column(column).expect("a column"));
1575                }
1576                gather.distinct()
1577            })
1578            .collect::<Vec<_>>();
1579
1580        for reversed in [false, true] {
1581            let split = path("capped");
1582            let mut writer = Writer::create(&split, "t", fields()).expect("a file");
1583            let preparer = writer.preparer();
1584            let mut prepared = runs()
1585                .into_iter()
1586                .map(|run| preparer.prepare(run).expect("prepared"))
1587                .collect::<Vec<_>>();
1588            for column in [0, 2] {
1589                assert!(preparer.coded.ceilings[column].load(Atomic::Relaxed) < u64::MAX);
1590            }
1591            if reversed {
1592                prepared.reverse();
1593            }
1594            for one in prepared {
1595                writer.append_prepared(one).expect("a stripe");
1596            }
1597            assert_eq!(estimates(&writer), want, "reversed {reversed}");
1598            writer.finish().expect("commit");
1599            fs::remove_file(split).expect("remove");
1600        }
1601    }
1602
1603    /// Two stripes merged in one order and written in the other read back as the rows they held,
1604    /// which is what two instances sharing a writer do whenever the second one's pages are built
1605    /// first.
1606    #[test]
1607    fn stripes_written_in_another_order_than_they_were_merged_read_back() {
1608        let path = path("crossed");
1609        let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1610        let preparer = writer.preparer();
1611        let mut merged = runs()
1612            .into_iter()
1613            .map(|run| writer.merge(preparer.prepare(run).expect("prepared")).expect("merged"))
1614            .map(|merged| merged.pages().expect("paged"))
1615            .collect::<Vec<_>>();
1616        merged.reverse();
1617        for paged in merged {
1618            writer.write(paged).expect("written");
1619        }
1620        writer.finish().expect("commit");
1621        check(&path);
1622        fs::remove_file(path).expect("remove");
1623    }
1624
1625    /// Stripes merged through a [`Merger`] write the same bytes as the writer merging them itself,
1626    /// and their rows are counted in when they are written.
1627    #[test]
1628    fn stripes_merged_through_a_merger_write_the_same_bytes_as_the_writer() {
1629        let alone = path("alone-merger");
1630        let mut writer = Writer::create(&alone, "t", fields()).expect("a file");
1631        for run in runs() {
1632            writer.append_stripe(run).expect("a stripe");
1633        }
1634        writer.finish().expect("commit");
1635
1636        let lent = path("lent");
1637        let mut writer = Writer::create(&lent, "t", fields()).expect("a file");
1638        let preparer = writer.preparer();
1639        let merger = writer.merger().expect("a merger");
1640        for run in runs() {
1641            let merged = merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1642            let mut paged = merged.pages().expect("paged");
1643            merger.give_back(&mut paged).expect("given back");
1644            writer.write(paged).expect("written");
1645        }
1646        assert_eq!(writer.table.rows, 13 * PART);
1647        writer.finish().expect("commit");
1648
1649        assert_eq!(fs::read(&alone).expect("read"), fs::read(&lent).expect("read"));
1650        check(&lent);
1651        fs::remove_file(alone).expect("remove");
1652        fs::remove_file(lent).expect("remove");
1653    }
1654
1655    /// A stripe fed a few parts at a time, with an empty batch and an empty chunk among them,
1656    /// writes the same bytes as the same stripe prepared whole.
1657    #[test]
1658    fn stripes_fed_in_batches_write_the_same_bytes_as_prepared_whole() {
1659        let whole = path("whole");
1660        let mut writer = Writer::create(&whole, "t", fields()).expect("a file");
1661        for run in runs() {
1662            writer.append_stripe(run).expect("a stripe");
1663        }
1664        writer.finish().expect("commit");
1665
1666        let fed = path("fed");
1667        let mut writer = Writer::create(&fed, "t", fields()).expect("a file");
1668        let preparer = writer.preparer();
1669        let merger = writer.merger().expect("a merger");
1670        let nothing = Chunk::new(
1671            fields()
1672                .iter()
1673                .map(|field| Vector::from_values(field.ty.clone(), &[]).expect("a column"))
1674                .collect(),
1675        )
1676        .expect("a chunk");
1677        for mut run in runs() {
1678            let mut building = preparer.start();
1679            preparer.feed(&mut building, Vec::new()).expect("fed nothing");
1680            while !run.is_empty() {
1681                let rest = run.split_off(2.min(run.len()));
1682                let mut batch = std::mem::replace(&mut run, rest);
1683                batch.push(((u64::MAX, 0), nothing.clone()));
1684                preparer.feed(&mut building, batch).expect("fed");
1685            }
1686            let merged =
1687                merger.merge(preparer.finish(building).expect("finished")).expect("merged");
1688            let mut paged = merged.pages().expect("paged");
1689            merger.give_back(&mut paged).expect("given back");
1690            writer.write(paged).expect("written");
1691        }
1692        writer.finish().expect("commit");
1693
1694        assert_eq!(fs::read(&whole).expect("read"), fs::read(&fed).expect("read"));
1695        check(&fed);
1696        fs::remove_file(whole).expect("remove");
1697        fs::remove_file(fed).expect("remove");
1698    }
1699
1700    /// A stripe started while `note` still had its dictionary, which the first stripe's merge then
1701    /// dropped, encodes the rest of `note` as pages and reads back.
1702    #[test]
1703    fn a_column_dropped_while_its_stripe_is_built_turns_to_pages() {
1704        let path = path("dropped-while-built");
1705        let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1706        let preparer = writer.preparer();
1707        let merger = writer.merger().expect("a merger");
1708        let write = |writer: &mut Writer, building: Building| {
1709            let merged =
1710                merger.merge(preparer.finish(building).expect("finished")).expect("merged");
1711            let mut paged = merged.pages().expect("paged");
1712            merger.give_back(&mut paged).expect("given back");
1713            writer.write(paged).expect("written");
1714        };
1715        let mut runs = runs().into_iter();
1716        let mut first = preparer.start();
1717        let mut second = preparer.start();
1718        let mut later = runs.next().expect("a run");
1719        preparer.feed(&mut second, later.drain(..2).collect()).expect("fed");
1720        preparer.feed(&mut first, runs.next().expect("a run")).expect("fed");
1721        write(&mut writer, first);
1722        let note = |building: &Building| {
1723            matches!(building.columns[2].lock().expect("unpoisoned").body, Body::Pages(..))
1724        };
1725        assert!(!note(&second), "still coded until it is fed again");
1726        preparer.feed(&mut second, later).expect("fed");
1727        assert!(note(&second), "turned to pages once fed after the drop");
1728        write(&mut writer, second);
1729        let mut last = preparer.start();
1730        preparer.feed(&mut last, runs.next().expect("a run")).expect("fed");
1731        write(&mut writer, last);
1732        writer.finish().expect("commit");
1733        check(&path);
1734        fs::remove_file(path).expect("remove");
1735    }
1736
1737    /// A stripe is held to [`STRIPE_PARTS`] across all its batches, not only within one.
1738    #[test]
1739    fn a_stripe_fed_more_parts_than_it_holds_is_refused() {
1740        let path = path("overfed");
1741        let writer = Writer::create(&path, "t", fields()).expect("a file");
1742        let preparer = writer.preparer();
1743        let mut building = preparer.start();
1744        preparer.feed(&mut building, stripe(0, STRIPE_PARTS - 1)).expect("fed");
1745        assert_eq!(building.parts(), STRIPE_PARTS - 1);
1746        assert!(preparer.feed(&mut building, stripe(STRIPE_PARTS, 2)).is_err());
1747        drop(writer);
1748        let _ = fs::remove_file(path);
1749    }
1750
1751    /// Stripes merged on several threads at once through one [`Merger`] and written in whatever
1752    /// order they finish read back as the rows they held.
1753    #[test]
1754    fn stripes_merged_on_several_threads_at_once_read_back() {
1755        let path = path("merged-at-once");
1756        let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1757        let preparer = writer.preparer();
1758        let merger = writer.merger().expect("a merger");
1759        let writer = Mutex::new(writer);
1760        std::thread::scope(|scope| {
1761            for run in runs() {
1762                let (preparer, merger, writer) = (&preparer, &merger, &writer);
1763                scope.spawn(move || {
1764                    let merged =
1765                        merger.merge(preparer.prepare(run).expect("prepared")).expect("merged");
1766                    let mut paged = merged.pages().expect("paged");
1767                    merger.give_back(&mut paged).expect("given back");
1768                    writer.lock().expect("the writer").write(paged).expect("written");
1769                });
1770            }
1771        });
1772        writer.into_inner().expect("the writer").finish().expect("commit");
1773        check(&path);
1774        fs::remove_file(path).expect("remove");
1775    }
1776
1777    /// A merge that comes after the table is closed is refused rather than merged into
1778    /// dictionaries nothing will write.
1779    #[test]
1780    fn a_merge_after_the_table_is_closed_is_refused() {
1781        let path = path("late");
1782        let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1783        let preparer = writer.preparer();
1784        let merger = writer.merger().expect("a merger");
1785        writer.finish().expect("commit");
1786        let prepared = preparer.prepare(stripe(0, 2)).expect("prepared");
1787        assert!(merger.merge(prepared).is_err());
1788        fs::remove_file(path).expect("remove");
1789    }
1790
1791    /// A chunk that is not the table's is refused when it reaches the writer, and the writer is not
1792    /// left counting its rows.
1793    #[test]
1794    fn a_stripe_of_another_table_is_refused_at_the_merge() {
1795        let path = path("refused");
1796        let mut writer = Writer::create(&path, "t", fields()).expect("a file");
1797        let other = Writer::create(path.with_extension("other"), "u", vec![fields().remove(0)])
1798            .expect("a file");
1799        let prepared = other.preparer().prepare(vec![]).expect("nothing to prepare");
1800        assert!(writer.merge(prepared).is_err());
1801        assert_eq!(writer.table.rows, 0);
1802        drop(other);
1803        fs::remove_file(path.with_extension("other")).expect("remove");
1804        fs::remove_file(path).expect("remove");
1805    }
1806
1807    /// A table whose `url` repeats twenty values for its first five parts and never repeats after,
1808    /// the way a log's URLs look once its first hours are past, next to a `city` that repeats
1809    /// throughout.
1810    fn turning(id: usize) -> [Value; 3] {
1811        let url = match id {
1812            _ if id % 13 == 0 => Value::Null,
1813            _ if id < 5 * PART => Value::Varchar(format!("https://example.com/{}", id % 20)),
1814            _ => Value::Varchar(format!("https://example.com/page/{id}")),
1815        };
1816        [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1817    }
1818
1819    fn turning_fields() -> Vec<Field> {
1820        vec![
1821            Field::required("id", LogicalType::BigInt),
1822            Field::new("city", LogicalType::Varchar),
1823            Field::new("url", LogicalType::Varchar),
1824        ]
1825    }
1826
1827    fn turning_stripe(
1828        rows: fn(usize) -> [Value; 3],
1829        first: usize,
1830        parts: usize,
1831    ) -> Vec<((u64, u64), Chunk)> {
1832        (first..first + parts)
1833            .map(|part| {
1834                let rows = (part * PART..(part + 1) * PART).map(rows).collect::<Vec<_>>();
1835                let column = |at: usize| {
1836                    let values = rows.iter().map(|row| row[at].clone()).collect::<Vec<_>>();
1837                    Vector::from_values(turning_fields()[at].ty.clone(), &values).expect("a column")
1838                };
1839                let chunk = Chunk::new(vec![column(0), column(1), column(2)]).expect("a chunk");
1840                ((part as u64, 0), chunk)
1841            })
1842            .collect()
1843    }
1844
1845    fn turning_runs() -> Vec<Vec<((u64, u64), Chunk)>> {
1846        vec![
1847            turning_stripe(turning, 0, 5),
1848            turning_stripe(turning, 5, 5),
1849            turning_stripe(turning, 10, 3),
1850        ]
1851    }
1852
1853    /// Reads every row of the turned table back and checks what the reader says about `url`.
1854    fn check_turned(path: &PathBuf) {
1855        let reader = Reader::open(path).expect("reopen");
1856        assert_eq!(reader.parts(), 13);
1857        assert_eq!(reader.table().demoted, [false, false, true], "only url is demoted");
1858        for part in 0..13 {
1859            let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1860            let url = chunk.column(2).expect("url");
1861            assert!(url.stable_dictionary_parts().is_none(), "part {part} hands out no codes");
1862            for at in 0..PART {
1863                let want = turning(part * PART + at);
1864                for (column, value) in want.iter().enumerate() {
1865                    assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1866                }
1867            }
1868        }
1869        // Everything the dictionary would have vouched for covers only the first stripes.
1870        assert_eq!(reader.distinct_values(2).expect("asked"), None);
1871        assert_eq!(reader.text_extremes(2).expect("asked"), None);
1872        assert_eq!(reader.exact_frequencies(2).expect("asked"), None);
1873        assert_eq!(reader.top_frequencies(2, 5).expect("asked"), None);
1874        assert!(!reader.skips_codes(0, 2, &[0]).expect("asked"), "no code proves a value absent");
1875        // `city` keeps its dictionary and all of it.
1876        assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1877        assert!(reader.text_extremes(1).expect("asked").is_some());
1878    }
1879
1880    /// A column whose second stripe is nearly all new values stops growing its dictionary there,
1881    /// whether the later stripes were prepared before that decision or after it, and every row
1882    /// reads back.
1883    #[test]
1884    fn a_column_that_turns_unique_is_demoted_and_reads_back() {
1885        let alone = path("demoted-alone");
1886        let mut writer = Writer::create(&alone, "t", turning_fields()).expect("a file");
1887        let preparer = writer.preparer();
1888        let mut runs = turning_runs().into_iter();
1889        writer.append_stripe(runs.next().expect("a run")).expect("a stripe");
1890        assert!(preparer.coded[2].load(Atomic::Relaxed), "url repeats in its first stripe");
1891        for run in runs {
1892            writer.append_stripe(run).expect("a stripe");
1893        }
1894        assert!(!preparer.coded[2].load(Atomic::Relaxed), "url was demoted");
1895        assert!(preparer.coded[1].load(Atomic::Relaxed), "city kept its dictionary");
1896        writer.finish().expect("commit");
1897        check_turned(&alone);
1898
1899        let split = path("demoted-split");
1900        let mut writer = Writer::create(&split, "t", turning_fields()).expect("a file");
1901        let preparer = writer.preparer();
1902        let prepared = turning_runs()
1903            .into_iter()
1904            .map(|run| preparer.prepare(run).expect("prepared"))
1905            .collect::<Vec<_>>();
1906        for one in prepared {
1907            writer.append_prepared(one).expect("a stripe");
1908        }
1909        writer.finish().expect("commit");
1910        check_turned(&split);
1911
1912        fs::remove_file(alone).expect("remove");
1913        fs::remove_file(split).expect("remove");
1914    }
1915
1916    /// A table where `url` takes a quarter of its rows as new values every stripe, which keeps its
1917    /// dictionary under the per-stripe rule, and `city` stops growing after its first stripe.
1918    fn growing(id: usize) -> [Value; 3] {
1919        let url = Value::Varchar(format!("https://example.com/{}", id / 4));
1920        [Value::BigInt(id as i64), Value::Varchar(format!("city {}", id % 13)), url]
1921    }
1922
1923    /// Once the dictionaries together pass the cap, the column that grew the most stops and the
1924    /// one that did not grow keeps its dictionary.
1925    #[test]
1926    fn the_dictionary_cap_demotes_the_column_that_grew_most() {
1927        let path = path("capped");
1928        let mut writer = Writer::create(&path, "t", turning_fields())
1929            .expect("a file")
1930            .with_dictionary_cap(1 << 30);
1931        let preparer = writer.preparer();
1932        writer.append_stripe(turning_stripe(growing, 0, 5)).expect("a stripe");
1933        assert!(preparer.coded[2].load(Atomic::Relaxed), "url is under the cap");
1934        assert!(preparer.coded[1].load(Atomic::Relaxed), "city is under the cap");
1935
1936        writer.coded.cap(1);
1937        writer.append_stripe(turning_stripe(growing, 5, 5)).expect("a stripe");
1938        assert!(!preparer.coded[2].load(Atomic::Relaxed), "url grew most and was demoted");
1939        assert!(preparer.coded[1].load(Atomic::Relaxed), "city grew nothing and keeps it");
1940        writer.append_stripe(turning_stripe(growing, 10, 3)).expect("a stripe");
1941        assert!(preparer.coded[1].load(Atomic::Relaxed), "city still grows nothing");
1942        writer.finish().expect("commit");
1943
1944        let reader = Reader::open(&path).expect("reopen");
1945        assert_eq!(reader.table().demoted, [false, false, true]);
1946        for part in 0..13 {
1947            let chunk = reader.read(part, &[0, 1, 2]).expect("a part");
1948            for at in 0..PART {
1949                let want = growing(part * PART + at);
1950                for (column, value) in want.iter().enumerate() {
1951                    assert_eq!(&chunk.value_at(at, column), value, "part {part} row {at}");
1952                }
1953            }
1954        }
1955        assert_eq!(reader.distinct_values(1).expect("asked"), Some(13));
1956        assert_eq!(reader.distinct_values(2).expect("asked"), None);
1957        fs::remove_file(path).expect("remove");
1958    }
1959}