gnitz-zset 0.1.1

The Z-set kernel of the gnitz database: schema, columnar batches, cursors and operators
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
//! [`MapPlan`] — the columnar driver behind every DBSP map and projection.
//!
//! A filter needs nothing on top of the resolved program and is a bare
//! `gnitz_expr::RowFilter` wherever one is held; a map is this plan, which adds
//! the PK, weight and column moves around the computed columns
//! `gnitz_expr::MapEval` writes.

use gnitz_expr::{ColCopy, ExprValidateErr, LogicalProgram, MapEval};

use super::reindex::{locate_key_col, FoldCols, ReindexPacker};
use crate::repr::{Batch, DirectWriter};
use crate::schema::{ColumnLocator, DerivedSchema, SchemaColumn, SchemaDescriptor, SchemaFacts, TypeCode};
use gnitz_wire::{zip_cells, FixedInt};

/// One map step's row window: source rows `[src, src + n)` onto destination rows
/// `[dst, dst + n)`. The three travel together through every columnar body
/// below, so they are one value rather than three positional `usize`s a caller
/// can transpose.
#[derive(Clone, Copy)]
struct RowWindow {
    src: usize,
    dst: usize,
    n: usize,
}

/// Average survivor-run length below which a computing map copies its survivors
/// into one range first: below it, the copy costs less than the kernel's setup
/// per range.
const COMPACT_RUN_LEN: usize = 128;

/// [`COMPACT_RUN_LEN`] for a copy-only map packing a reindex key, whose
/// per-range setup is smaller.
const PACK_COMPACT_RUN_LEN: usize = 24;

/// Where a map's output PK region comes from. Owned by the plan rather than
/// passed per call, so the region cannot be left unwritten between two
/// statements and no caller can pair a plan with the wrong stamp.
enum PkSource {
    /// Copy the input PK region verbatim, into an output schema whose PK is the
    /// input's.
    Inherit,
    /// Pack the reindex columns' OPK bytes contiguously into the output PK — the
    /// `_join_pk` of an equijoin / GROUP BY repartition. The output stride
    /// legitimately differs from the input's (U64 input → U128 synthetic PK).
    Pack(ReindexPacker),
    /// Hash each output row's payload columns into its PK.
    HashRow(FoldCols),
}

// ---------------------------------------------------------------------------
// Helpers
// ---------------------------------------------------------------------------

/// Copy a single column from `in_batch` to `output` over `w`, a German-string
/// cell verbatim: the caller rebases the column. The source is read at its own
/// type's width and sign/zero-extended into a wider (promoted) destination slot.
fn copy_column(
    in_batch: &Batch,
    output: &mut DirectWriter<'_>,
    &ColCopy {
        src: src_loc,
        slot: dst_payload,
        width: stride,
    }: &ColCopy,
    w: RowWindow,
) {
    let RowWindow { src: src_start, dst: dst_base, n } = w;
    // Destructured once, so each arm below is one column kernel over the window.
    match src_loc {
        ColumnLocator::Pk { byte_off, size, type_code } => {
            // The PK region holds OPK bytes, decoded back to native LE at the
            // source column's own width.
            let pk_stride = in_batch.schema().pk_stride();
            let pk = &in_batch.pk_data()[src_start * pk_stride..(src_start + n) * pk_stride];
            let dst = &mut output.col_mut(dst_payload)[dst_base * stride..(dst_base + n) * stride];
            let (off, src_stride) = (byte_off as usize, size as usize);
            if src_stride == stride {
                gnitz_wire::decode_pk_cells(pk, pk_stride, off, stride, type_code.is_signed_int(), dst);
            } else {
                widen_column(pk, pk_stride, off, type_code, true, stride, dst);
            }
        }
        ColumnLocator::Payload { slot, size, type_code } => {
            let in_pi = slot as usize;
            let src_stride = size as usize; // source read width
            debug_assert!(
                !type_code.is_german_string() || (src_stride, stride) == (16, 16),
                "German-string column moved at a non-16-byte stride",
            );
            debug_assert!(
                (src_start + n) * src_stride <= in_batch.col_data(in_pi).len(),
                "copy_column: source column {in_pi} is shorter than rows [{src_start}, {}) at stride {src_stride}",
                src_start + n
            );
            let src = &in_batch.col_data(in_pi)[src_start * src_stride..(src_start + n) * src_stride];
            let dst = &mut output.col_mut(dst_payload)[dst_base * stride..(dst_base + n) * stride];
            if src_stride == stride {
                dst.copy_from_slice(src);
            } else {
                widen_column(src, src_stride, 0, type_code, false, stride, dst);
            }
        }
    }
}

/// [`widen_cells`] at the slot width `dw`: 2, 4 or 8, a widened slot being a fixed int.
fn widen_column(src: &[u8], src_stride: usize, off: usize, tc: TypeCode, pk: bool, dw: usize, dst: &mut [u8]) {
    let fi = FixedInt::from_type_code(tc).expect("a widened column is a fixed int");
    match dw {
        2 => widen_cells::<2>(src, src_stride, off, fi, pk, dst),
        4 => widen_cells::<4>(src, src_stride, off, fi, pk, dst),
        8 => widen_cells::<8>(src, src_stride, off, fi, pk, dst),
        other => unreachable!("a widened slot is 2/4/8 bytes, got {other}"),
    }
}

/// Widen each `fi` cell at byte `off` of `src`'s `src_stride`-byte rows — a PK column's
/// OPK cell when `pk`, else a native one — into `dst`'s `DW`-byte native cells.
fn widen_cells<const DW: usize>(src: &[u8], src_stride: usize, off: usize, fi: FixedInt, pk: bool, dst: &mut [u8]) {
    let out = dst.as_chunks_mut::<DW>().0.iter_mut();
    gnitz_wire::for_each_fixed_int!(fi, |FI| {
        const W: usize = FI.width();
        let store = |v: i64, d: &mut [u8; DW]| *d = v.to_le_bytes()[..DW].try_into().unwrap();
        match pk {
            true => zip_cells::<W, _>(src, src_stride, off, out, |c, d| {
                store(gnitz_wire::decode_opk_i64(c, FI), d)
            }),
            false => zip_cells::<W, _>(src, src_stride, off, out, |c, d| store(FI.decode_le_i64(c), d)),
        }
    });
}

// ---------------------------------------------------------------------------
// MapPlan
// ---------------------------------------------------------------------------

/// Map/projection: the resolved program's columnar moves and computed columns,
/// the null permutation, where the output PK comes from, and the output schema
/// it stamps.
pub struct MapPlan {
    ev: MapEval,
    /// Where the output PK region comes from.
    pk_source: PkSource,
    /// The input's German-string payload slots some copy reads.
    copied_string_slots: u64,
    /// A row with any of these null bits set is dropped: a
    /// [`gnitz_wire::NullKeys::Drop`] reindex's nullable key columns.
    null_key_mask: u64,
    out_schema: SchemaDescriptor,
}

/// Set the PK of each of `output`'s rows, whose payload and null words are
/// written, to the [`FoldCols`] digest of its payload columns, so equal rows
/// share a PK. Two distinct rows whose 128-bit digests collide become one
/// element; accepted, not checked.
fn reindex_hash_row(output: &mut DirectWriter<'_>, fold: &FoldCols) {
    let (first, n) = (output.first_row(), output.rows());
    // The PK *is* the digest, so the OPK region is its big-endian bytes.
    const KEY_BYTES: usize = std::mem::size_of::<u128>();
    assert_eq!(output.schema.pk_stride(), KEY_BYTES, "a hash-row PK is one U128 column");
    // Hashing borrows `output` and the write-back mutates it, so keys are staged
    // per chunk.
    const CHUNK: usize = 256;
    let mut keys = [0u128; CHUNK];
    let mut start = 0;
    while start < n {
        let end = (start + CHUNK).min(n);
        {
            let mb = output.written();
            for (key, row) in keys.iter_mut().zip(first + start..first + end) {
                *key = fold.key_row(&mb, row, mb.get_null_word(row));
            }
        }
        let pk = &mut output.pk_mut()[start * KEY_BYTES..end * KEY_BYTES];
        for (key, dst) in keys.iter().zip(pk.as_chunks_mut::<KEY_BYTES>().0) {
            *dst = key.to_be_bytes();
        }
        start = end;
    }
}

/// Output schema of a computed-projection `Map`: the input's PK region (the map
/// inherits it verbatim, [`PkSource::Inherit`]), then one payload column per
/// declared `(type_code, nullable)` slot. The declared slots ARE the layout;
/// `MapPlan::from_map` is what catches a program disagreeing with them.
fn compute_map_output_schema(
    in_schema: &SchemaDescriptor,
    out_cols: &[(TypeCode, bool)],
) -> Result<SchemaDescriptor, String> {
    let mut b = DerivedSchema::new();
    b.push_pk_of(in_schema);
    for &(tc, nullable) in out_cols {
        b.push(SchemaColumn::new(tc, nullable));
    }
    b.finish().map_err(|e| format!("compute map: output {e}"))
}

/// Output schema of a HashRow Map: a U128 PK, then each projected column at its
/// slot type and its source nullability. Typed at the slot, the promotion is
/// what the copy kernel screens.
fn hashrow_output_schema(
    in_schema: &SchemaDescriptor,
    cols: &[gnitz_wire::ReindexSlot],
) -> Result<SchemaDescriptor, String> {
    let mut b = DerivedSchema::new();
    b.push_pk(SchemaColumn::new(crate::schema::TypeCode::U128, false));
    for &(c, t) in cols {
        // A key column, not merely an in-range one — the screen the reindex and
        // top-N key kinds clear at this same boundary.
        locate_key_col(in_schema, c, "hash-row map")?;
        let src = in_schema.columns[c as usize];
        b.push(SchemaColumn::new(t, src.nullable));
    }
    b.finish().map_err(|e| format!("hash-row map: output {e}"))
}

impl MapPlan {
    /// The plan for a circuit's MAP node: one derivation of its `(output schema,
    /// map program, PK source)` triple, and the trust boundary each kind's
    /// client-supplied column list clears. Every kind ends in the same
    /// [`Self::from_map`], so an elided map is validated like any other.
    pub fn from_wire(in_schema: &SchemaDescriptor, mk: &gnitz_wire::MapKind) -> Result<Self, String> {
        let mut null_key_mask = 0;
        let (out_schema, prog, pk_source) = match mk {
            gnitz_wire::MapKind::Compute(map) => return Self::from_compute_map(in_schema, map),

            gnitz_wire::MapKind::Reindex { keep, key, nulls, .. } => {
                let packer = ReindexPacker::new(in_schema, key)?;
                let out_schema = packer.output_schema(in_schema, keep)?;
                if *nulls == gnitz_wire::NullKeys::Drop {
                    null_key_mask = key
                        .iter()
                        .filter_map(|&(c, _)| in_schema.payload_slot(c as usize))
                        .fold(0, |mask, slot| mask | 1u64 << slot)
                        & in_schema.nullable_payload_slots();
                }
                (out_schema, LogicalProgram::copy_cols(keep), PkSource::Pack(packer))
            }

            gnitz_wire::MapKind::HashRow { cols } => {
                let out_schema = hashrow_output_schema(in_schema, cols)?;
                let proj: Vec<u32> = cols.iter().map(|&(c, _)| c).collect();
                let fold = FoldCols::new(out_schema.payload_locators());
                (out_schema, LogicalProgram::copy_cols(&proj), PkSource::HashRow(fold))
            }

            gnitz_wire::MapKind::Projection(cols) => {
                let out_schema =
                    crate::schema::project_schema(in_schema, cols).map_err(|e| format!("projection map: {e}"))?;
                (out_schema, LogicalProgram::copy_cols(cols), PkSource::Inherit)
            }
        };
        let plan = Self::from_map(prog, in_schema, &out_schema, pk_source)
            .map_err(|e| format!("map: program/schema mismatch: {e}"))?;
        Ok(MapPlan { null_key_mask, ..plan })
    }

    /// A computed projection: [`gnitz_wire::MapKind::Compute`]'s whole body, and
    /// also a read spec's pre-map. The output schema is derived, never shipped —
    /// the map inherits the input's PK region verbatim ([`PkSource::Inherit`]),
    /// so no caller can describe a PK region the map does not produce.
    pub(crate) fn from_compute_map(in_schema: &SchemaDescriptor, map: &gnitz_wire::ComputeMap) -> Result<Self, String> {
        let out_schema = compute_map_output_schema(in_schema, &map.out_cols)?;
        // The only map whose program is client bytes; every other kind builds
        // one from a column list. Rejected, not skipped: skipping a corrupt blob
        // would leave the output at the default empty schema.
        let prog = LogicalProgram::from_blob(&map.program).map_err(|e| format!("map: invalid program: {e}"))?;
        Self::from_map(prog, in_schema, &out_schema, PkSource::Inherit)
            .map_err(|e| format!("map: program/schema mismatch: {e}"))
    }

    /// Map plan from a logical expression program. A pure projection is the
    /// special case where the program computes nothing and every sink is a
    /// column copy (see [`LogicalProgram::copy_cols`]): the plan reduces to the
    /// copy list and the resolved program's null permutation.
    fn from_map(
        logical: LogicalProgram,
        in_schema: &SchemaDescriptor,
        out_schema: &SchemaDescriptor,
        pk_source: PkSource,
    ) -> Result<Self, ExprValidateErr> {
        let ev = logical.resolve_map(in_schema, out_schema)?;
        let copied_string_slots = ev.copies().iter().fold(0u64, |slots, c| match c.src {
            ColumnLocator::Payload { slot, type_code, .. } if type_code.is_german_string() => slots | 1u64 << slot,
            _ => slots,
        });

        Ok(MapPlan {
            ev,
            pk_source,
            copied_string_slots,
            null_key_mask: 0,
            out_schema: *out_schema,
        })
    }

    /// The schema this plan stamps on its output.
    pub fn out_schema(&self) -> &SchemaDescriptor {
        &self.out_schema
    }

    /// The input column behind each output payload slot, when this map only
    /// re-keys its rows onto leading bytes of their own PK: the input, read in
    /// its own order, then holds each output key's rows together. `None` for
    /// every other map.
    pub fn rekeys_onto_pk_prefix(&self) -> Option<Vec<ColumnLocator>> {
        let PkSource::Pack(packer) = &self.pk_source else {
            return None;
        };
        packer.pk_range().filter(|&(at, _)| at == 0)?;
        // A PK column is never NULL, and a reindex copies each kept column at
        // its own type.
        debug_assert!(self.null_key_mask == 0 && !self.ev.emits_anything());
        debug_assert!(self.ev.copies().iter().all(|c| c.width == c.src.size()));
        Some(self.ev.copies().iter().map(|c| c.src).collect())
    }

    /// Whether some row can be dropped: a [`gnitz_wire::NullKeys::Drop`] re-key over
    /// a nullable key column.
    pub fn drops_null_keys(&self) -> bool {
        self.null_key_mask != 0
    }

    /// True iff running this map would reproduce its input batch. A compiler
    /// elides such a node entirely and lets its consumers read the input.
    pub fn is_identity(&self) -> bool {
        matches!(self.pk_source, PkSource::Inherit) && self.ev.is_identity()
    }

    /// The map over `in_batch`, less a [`gnitz_wire::NullKeys::Drop`] reindex's
    /// NULL-keyed rows, into a fresh unconsolidated output.
    pub fn evaluate_map_batch(&mut self, in_batch: &Batch) -> Batch {
        let whole = [(0, in_batch.count)];
        let runs;
        let ranges: &[(usize, usize)] = match self.null_key_mask {
            0 => &whole,
            mask => {
                runs = in_batch.runs_without_nulls(mask);
                &runs
            }
        };
        let n = crate::repr::range_rows(ranges);
        if n == 0 {
            return Batch::empty_with_schema(&self.out_schema);
        }
        // Uninitialized: `validate` makes every map write every payload slot,
        // and the calls below cover the PK, weight and null regions.
        let mut output = Batch::with_capacity(&self.out_schema, n);
        self.append_map_ranges(in_batch, &mut output, ranges);
        output
    }

    /// Map every `[start, end)` range of `src`, in order, onto the tail of `out`,
    /// which is in this plan's output schema. Ranges too short to feed the
    /// per-window kernel are gathered into one first.
    pub(crate) fn append_map_ranges(&mut self, src: &Batch, out: &mut Batch, ranges: &[(usize, usize)]) {
        let total = crate::repr::range_rows(ranges);
        if total == 0 {
            return;
        }
        let compact_below = match (self.ev.emits_anything(), &self.pk_source) {
            (true, _) => COMPACT_RUN_LEN,
            (false, PkSource::Pack(_)) => PACK_COMPACT_RUN_LEN,
            (false, _) => 0,
        };
        let starves_kernel = ranges.len() > 1 && total < ranges.len() * compact_below;
        let compacted = starves_kernel.then(|| Batch::from_ranges(src, ranges, 0));
        let whole = [(0, total)];
        let (src, ranges) = match &compacted {
            Some(c) => (c, &whole[..]),
            None => (src, ranges),
        };
        let heap_at = out.carry_heap(&src.as_mem_batch(), self.copied_string_slots, ranges);
        // For relocated copies and string emits alike.
        if heap_at.is_none() {
            out.reserve_blob(crate::repr::prorated_blob_cap(src.blob().len(), src.count, total));
        }
        out.append_session(total).write(total, |out| {
            let mut dst = 0;
            for &(start, end) in ranges {
                let w = RowWindow { src: start, dst, n: end - start };
                self.map_rows_into(src, out, w);
                dst += w.n;
            }
            for c in self.ev.copies() {
                if matches!(c.src, ColumnLocator::Payload { type_code, .. } if type_code.is_german_string()) {
                    out.rebase_string_col(c.slot, src.blob(), heap_at);
                }
            }
            // The one source that keys on the finished output row.
            if let PkSource::HashRow(fold) = &self.pk_source {
                reindex_hash_row(out, fold);
            }
        });
    }

    /// Map one row window of `output`'s rows: PK, weight, column moves, then the
    /// null words and computed columns.
    ///
    /// `#[inline(always)]`: it runs once per window, and a window can be one row.
    #[inline(always)]
    fn map_rows_into(&mut self, in_batch: &Batch, output: &mut DirectWriter<'_>, w: RowWindow) {
        let RowWindow { src: src_start, dst: dst_base, n } = w;
        // Both PK sources that read the *input* row, so both belong to the
        // window rather than to a pass over the finished batch.
        match &self.pk_source {
            PkSource::Inherit => {
                let pk_st = in_batch.schema().pk_stride();
                debug_assert_eq!(
                    pk_st,
                    output.schema.pk_stride(),
                    "PkSource::Inherit: PK stride mismatch"
                );
                output.pk_mut()[dst_base * pk_st..(dst_base + n) * pk_st]
                    .copy_from_slice(&in_batch.pk_data()[src_start * pk_st..(src_start + n) * pk_st]);
            }
            PkSource::Pack(packer) => {
                debug_assert_eq!(output.schema.pk_stride(), packer.out_stride);
                let stride = packer.out_stride;
                let pk = &mut output.pk_mut()[dst_base * stride..];
                packer.pack_rows(pk, stride, &in_batch.as_mem_batch(), &[(src_start, src_start + n)]);
            }
            // Hashes the finished output row, so `append_map_ranges` stamps it
            // once every window's payload is written.
            PkSource::HashRow(_) => {}
        }
        output.weight_mut()[dst_base * 8..(dst_base + n) * 8]
            .copy_from_slice(&in_batch.weight_data()[src_start * 8..(src_start + n) * 8]);

        for c in self.ev.copies() {
            copy_column(in_batch, output, c, w);
        }
        self.ev
            .write_computed(&in_batch.as_mem_batch(), src_start, n, output, dst_base);
    }
}

// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------

#[cfg(test)]
#[path = "tests/map.rs"]
mod tests;

#[cfg(test)]
#[path = "benches/map.rs"]
mod bench;