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