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