1use std::sync::Arc;
6
7use nodedb_codec::{ColumnCodec, ColumnTypeHint, ResolvedColumnCodec};
8use nodedb_mem::{EngineId, MemoryGovernor};
9use nodedb_types::columnar::{ColumnType, ColumnarSchema};
10
11use crate::error::ColumnarError;
12use crate::format::{ColumnMeta, HEADER_SIZE, SegmentFooter, SegmentHeader};
13use crate::memtable::ColumnData;
14
15use super::block::encode_column_blocks;
16use super::encode::compute_schema_hash;
17
18pub const PROFILE_PLAIN: u8 = 0;
20pub const PROFILE_TIMESERIES: u8 = 1;
21pub const PROFILE_SPATIAL: u8 = 2;
22
23pub struct SegmentWriter {
29 profile_tag: u8,
30 governor: Option<Arc<MemoryGovernor>>,
33}
34
35impl SegmentWriter {
36 pub fn new(profile_tag: u8) -> Self {
38 Self {
39 profile_tag,
40 governor: None,
41 }
42 }
43
44 pub fn with_governor(profile_tag: u8, governor: Arc<MemoryGovernor>) -> Self {
46 Self {
47 profile_tag,
48 governor: Some(governor),
49 }
50 }
51
52 pub fn plain() -> Self {
54 Self::new(PROFILE_PLAIN)
55 }
56
57 pub fn write_segment(
66 &self,
67 schema: &ColumnarSchema,
68 columns: &[ColumnData],
69 row_count: usize,
70 kek: Option<&nodedb_wal::crypto::WalEncryptionKey>,
71 ) -> Result<Vec<u8>, ColumnarError> {
72 if row_count == 0 {
73 return Err(ColumnarError::EmptyMemtable);
74 }
75 if columns.len() != schema.columns.len() {
76 return Err(ColumnarError::SchemaMismatch {
77 expected: schema.columns.len(),
78 got: columns.len(),
79 });
80 }
81
82 let mut buf = Vec::new();
83
84 buf.extend_from_slice(&SegmentHeader::current().to_bytes());
86
87 let _metas_guard = self
89 .governor
90 .as_ref()
91 .map(|g| {
92 g.reserve(
93 EngineId::Columnar,
94 columns.len() * std::mem::size_of::<ColumnMeta>(),
95 )
96 })
97 .transpose()?;
98 let mut column_metas = Vec::with_capacity(columns.len());
99
100 for (i, (col_def, col_data)) in schema.columns.iter().zip(columns.iter()).enumerate() {
101 let col_start = buf.len() as u64;
102
103 let codec = select_codec_for_profile(&col_def.column_type, self.profile_tag);
105
106 let block_stats = encode_column_blocks(
108 &mut buf,
109 col_data,
110 &col_def.column_type,
111 codec,
112 row_count,
113 self.governor.as_ref(),
114 )?;
115
116 let col_end = buf.len() as u64;
117
118 let (effective_codec, dictionary) = match col_data {
121 ColumnData::DictEncoded { dictionary, .. } => (
122 ResolvedColumnCodec::DeltaFastLanesLz4,
123 Some(dictionary.clone()),
124 ),
125 _ => (codec, None),
126 };
127
128 column_metas.push(ColumnMeta {
129 name: col_def.name.clone(),
130 offset: col_start - HEADER_SIZE as u64,
131 length: col_end - col_start,
132 codec: effective_codec,
133 block_count: block_stats.len() as u32,
134 block_stats,
135 dictionary,
136 });
137
138 let _ = i; }
140
141 let schema_hash = compute_schema_hash(schema);
143
144 let footer = SegmentFooter {
146 schema_hash,
147 column_count: schema.columns.len() as u32,
148 row_count: row_count as u64,
149 profile_tag: self.profile_tag,
150 columns: column_metas,
151 };
152 let footer_bytes = footer.to_bytes()?;
153 buf.extend_from_slice(&footer_bytes);
154
155 if let Some(key) = kek {
157 return crate::encrypt::encrypt_segment(key, &buf);
158 }
159
160 Ok(buf)
161 }
162}
163
164pub fn select_codec_for_profile(col_type: &ColumnType, profile_tag: u8) -> ResolvedColumnCodec {
172 if profile_tag == PROFILE_TIMESERIES && matches!(col_type, ColumnType::Float64) {
174 return ResolvedColumnCodec::Gorilla;
175 }
176 if profile_tag == PROFILE_TIMESERIES
178 && matches!(col_type, ColumnType::Timestamp | ColumnType::Timestamptz)
179 {
180 return ResolvedColumnCodec::DeltaFastLanesLz4;
181 }
182 select_codec(col_type)
183}
184
185fn select_codec(col_type: &ColumnType) -> ResolvedColumnCodec {
190 let hint = match col_type {
191 ColumnType::Int64 => ColumnTypeHint::Int64,
192 ColumnType::Float64 => ColumnTypeHint::Float64,
193 ColumnType::Timestamp | ColumnType::Timestamptz | ColumnType::SystemTimestamp => {
194 ColumnTypeHint::Timestamp
195 }
196 ColumnType::String
197 | ColumnType::Geometry
198 | ColumnType::Regex
199 | ColumnType::SparseVector => ColumnTypeHint::String,
200 ColumnType::Bool
201 | ColumnType::Bytes
202 | ColumnType::Decimal { .. }
203 | ColumnType::Uuid
204 | ColumnType::Ulid => {
205 return ResolvedColumnCodec::Lz4;
206 }
207 ColumnType::Json
208 | ColumnType::Array
209 | ColumnType::Set
210 | ColumnType::Range
211 | ColumnType::Record => {
212 return ResolvedColumnCodec::Lz4;
213 }
214 ColumnType::Duration => ColumnTypeHint::Int64, ColumnType::Vector(_) => {
216 return ResolvedColumnCodec::Lz4;
217 }
218 _ => {
221 return ResolvedColumnCodec::Lz4;
222 }
223 };
224 nodedb_codec::detect_codec(ColumnCodec::Auto, hint)
228 .try_resolve()
229 .unwrap_or(ResolvedColumnCodec::Lz4)
230}
231
232#[cfg(test)]
233mod tests {
234 use nodedb_types::columnar::{ColumnDef, ColumnType, ColumnarSchema};
235 use nodedb_types::value::Value;
236
237 use super::*;
238 use crate::format::{SegmentFooter, SegmentHeader};
239 use crate::memtable::ColumnarMemtable;
240
241 fn analytics_schema() -> ColumnarSchema {
242 ColumnarSchema::new(vec![
243 ColumnDef::required("id", ColumnType::Int64).with_primary_key(),
244 ColumnDef::required("name", ColumnType::String),
245 ColumnDef::nullable("score", ColumnType::Float64),
246 ])
247 .expect("valid")
248 }
249
250 #[test]
255 fn auto_codec_resolves_to_concrete_before_write() {
256 let schema = ColumnarSchema::new(vec![
257 ColumnDef::required("id", ColumnType::Int64).with_primary_key(),
258 ColumnDef::required("name", ColumnType::String),
259 ColumnDef::nullable("score", ColumnType::Float64),
260 ])
261 .expect("valid");
262
263 let mut mt = ColumnarMemtable::new(&schema);
264 for i in 0..10 {
265 mt.append_row(&[
266 Value::Integer(i),
267 Value::String(format!("item_{i}")),
268 Value::Float(i as f64 * 1.5),
269 ])
270 .expect("append");
271 }
272 let (schema, columns, row_count) = mt.drain();
273 let writer = SegmentWriter::plain();
274 let segment = writer
275 .write_segment(&schema, &columns, row_count, None)
276 .expect("write must succeed");
277
278 let footer = SegmentFooter::from_segment_tail(&segment).expect("valid footer");
279
280 for col in &footer.columns {
283 let encoded = zerompk::to_msgpack_vec(&col.codec).expect("serialize");
287 let discriminant_byte = *encoded.last().expect("non-empty");
290 assert_ne!(
291 discriminant_byte, 0,
292 "column '{}' has Auto discriminant (0) on disk — resolve was skipped",
293 col.name
294 );
295 }
296 }
297
298 #[test]
300 fn auto_codec_int64_resolves_to_non_raw() {
301 use nodedb_codec::ResolvedColumnCodec;
302
303 let schema = ColumnarSchema::new(vec![
304 ColumnDef::required("val", ColumnType::Int64).with_primary_key(),
305 ])
306 .expect("valid");
307
308 let mut mt = ColumnarMemtable::new(&schema);
309 for i in 0..10 {
310 mt.append_row(&[Value::Integer(i)]).expect("append");
311 }
312 let (schema, columns, row_count) = mt.drain();
313 let writer = SegmentWriter::plain();
314 let segment = writer
315 .write_segment(&schema, &columns, row_count, None)
316 .expect("write");
317 let footer = SegmentFooter::from_segment_tail(&segment).expect("footer");
318
319 let codec = footer.columns[0].codec;
321 assert_ne!(
322 codec,
323 ResolvedColumnCodec::Raw,
324 "Int64 should not resolve to Raw"
325 );
326 }
327
328 #[test]
329 fn write_segment_roundtrip() {
330 let schema = analytics_schema();
331 let mut mt = ColumnarMemtable::new(&schema);
332
333 for i in 0..100 {
334 mt.append_row(&[
335 Value::Integer(i),
336 Value::String(format!("user_{i}")),
337 if i % 3 == 0 {
338 Value::Null
339 } else {
340 Value::Float(i as f64 * 0.25)
341 },
342 ])
343 .expect("append");
344 }
345
346 let (schema, columns, row_count) = mt.drain();
347 let writer = SegmentWriter::plain();
348 let segment = writer
349 .write_segment(&schema, &columns, row_count, None)
350 .expect("write");
351
352 let header = SegmentHeader::from_bytes(&segment).expect("valid header");
354 assert_eq!(header.magic, *b"NDBS");
355 assert_eq!(header.version_major, 1);
356
357 let footer = SegmentFooter::from_segment_tail(&segment).expect("valid footer");
359 assert_eq!(footer.column_count, 3);
360 assert_eq!(footer.row_count, 100);
361 assert_eq!(footer.profile_tag, PROFILE_PLAIN);
362 assert_eq!(footer.columns.len(), 3);
363
364 assert_eq!(footer.columns[0].name, "id");
366 assert_eq!(footer.columns[1].name, "name");
367 assert_eq!(footer.columns[2].name, "score");
368
369 assert_eq!(footer.columns[0].block_count, 1);
371 assert_eq!(footer.columns[0].block_stats[0].row_count, 100);
372
373 assert_eq!(footer.columns[0].block_stats[0].min, 0.0);
375 assert_eq!(footer.columns[0].block_stats[0].max, 99.0);
376 assert_eq!(footer.columns[0].block_stats[0].null_count, 0);
377
378 assert_eq!(footer.columns[2].block_stats[0].null_count, 34);
380 }
381
382 #[test]
383 fn write_segment_multi_block() {
384 let schema =
385 ColumnarSchema::new(vec![ColumnDef::required("x", ColumnType::Int64)]).expect("valid");
386
387 let mut mt = ColumnarMemtable::new(&schema);
388 for i in 0..2500 {
389 mt.append_row(&[Value::Integer(i)]).expect("append");
390 }
391
392 let (schema, columns, row_count) = mt.drain();
393 let writer = SegmentWriter::plain();
394 let segment = writer
395 .write_segment(&schema, &columns, row_count, None)
396 .expect("write");
397
398 let footer = SegmentFooter::from_segment_tail(&segment).expect("valid footer");
399 assert_eq!(footer.row_count, 2500);
400
401 assert_eq!(footer.columns[0].block_count, 3);
403 assert_eq!(footer.columns[0].block_stats[0].row_count, 1024);
404 assert_eq!(footer.columns[0].block_stats[1].row_count, 1024);
405 assert_eq!(footer.columns[0].block_stats[2].row_count, 452);
406
407 assert_eq!(footer.columns[0].block_stats[0].min, 0.0);
409 assert_eq!(footer.columns[0].block_stats[0].max, 1023.0);
410 assert_eq!(footer.columns[0].block_stats[2].min, 2048.0);
412 assert_eq!(footer.columns[0].block_stats[2].max, 2499.0);
413 }
414
415 #[test]
416 fn write_segment_empty_rejected() {
417 let schema = analytics_schema();
418 let mt = ColumnarMemtable::new(&schema);
419 let (schema, columns, row_count) = {
420 let mut m = mt;
421 m.drain()
422 };
423 let writer = SegmentWriter::plain();
424 assert!(matches!(
425 writer.write_segment(&schema, &columns, row_count, None),
426 Err(ColumnarError::EmptyMemtable)
427 ));
428 }
429
430 #[test]
431 fn block_stats_predicate_pushdown() {
432 let schema = analytics_schema();
433 let mut mt = ColumnarMemtable::new(&schema);
434
435 for i in 0..50 {
436 mt.append_row(&[
437 Value::Integer(i + 100),
438 Value::String(format!("item_{i}")),
439 Value::Float(i as f64 + 10.0),
440 ])
441 .expect("append");
442 }
443
444 let (schema, columns, row_count) = mt.drain();
445 let writer = SegmentWriter::plain();
446 let segment = writer
447 .write_segment(&schema, &columns, row_count, None)
448 .expect("write");
449 let footer = SegmentFooter::from_segment_tail(&segment).expect("valid");
450
451 use crate::predicate::ScanPredicate;
452
453 let id_stats = &footer.columns[0].block_stats[0];
454 assert!(ScanPredicate::gt(0, 200.0).can_skip_block(id_stats)); assert!(!ScanPredicate::gt(0, 120.0).can_skip_block(id_stats)); assert!(ScanPredicate::lt(0, 50.0).can_skip_block(id_stats)); assert!(ScanPredicate::eq(0, 200.0).can_skip_block(id_stats)); assert!(!ScanPredicate::eq(0, 125.0).can_skip_block(id_stats)); }
461
462 #[test]
463 fn string_block_stats_zone_map() {
464 let schema = ColumnarSchema::new(vec![ColumnDef::required("tag", ColumnType::String)])
466 .expect("valid");
467
468 let mut mt = ColumnarMemtable::new(&schema);
469 let values: Vec<String> = (0..20).map(|i| format!("item_{i:02}")).collect();
472 for name in &values {
473 mt.append_row(&[Value::String(name.clone())])
474 .expect("append");
475 }
476 mt.append_row(&[Value::String("apple".into())])
478 .expect("append");
479 mt.append_row(&[Value::String("date".into())])
480 .expect("append");
481
482 let (schema, columns, row_count) = mt.drain();
483 let writer = SegmentWriter::plain();
484 let segment = writer
485 .write_segment(&schema, &columns, row_count, None)
486 .expect("write");
487 let footer = SegmentFooter::from_segment_tail(&segment).expect("footer");
488
489 let stats = &footer.columns[0].block_stats[0];
490 assert!(stats.str_min.is_some(), "str_min should be populated");
491 assert!(stats.str_max.is_some(), "str_max should be populated");
492 assert_eq!(stats.str_min.as_deref(), Some("apple"));
494 assert_eq!(stats.str_max.as_deref(), Some("item_19"));
495
496 assert!(
498 stats.bloom.is_some(),
499 "bloom should be populated for >16 distinct values"
500 );
501
502 use crate::predicate::ScanPredicate;
503
504 assert!(ScanPredicate::str_eq(0, "aaa").can_skip_block(stats));
506 assert!(ScanPredicate::str_eq(0, "zzz").can_skip_block(stats));
508 assert!(!ScanPredicate::str_eq(0, "date").can_skip_block(stats));
510 assert!(ScanPredicate::str_gt(0, "item_19").can_skip_block(stats));
512 assert!(ScanPredicate::str_lt(0, "apple").can_skip_block(stats));
514 }
515
516 #[test]
520 fn timestamp_large_value_roundtrip() {
521 use crate::predicate::ScanPredicate;
522
523 let schema = ColumnarSchema::new(vec![
524 ColumnDef::required("ts", ColumnType::Timestamp).with_primary_key(),
525 ])
526 .expect("valid schema");
527
528 let base: i64 = 10_413_792_000_000_000;
532 let target = base + 500_000; let mut mt = ColumnarMemtable::new(&schema);
535 for delta in 0..10i64 {
536 mt.append_row(&[Value::Integer(base + delta * 100_000)])
537 .expect("append");
538 }
539
540 let (schema, columns, row_count) = mt.drain();
541 let segment = SegmentWriter::plain()
542 .write_segment(&schema, &columns, row_count, None)
543 .expect("write");
544 let footer = SegmentFooter::from_segment_tail(&segment).expect("footer");
545
546 let stats = &footer.columns[0].block_stats[0];
547
548 assert!(
550 stats.min_i64.is_some(),
551 "min_i64 must be set for timestamp columns"
552 );
553 assert!(
554 stats.max_i64.is_some(),
555 "max_i64 must be set for timestamp columns"
556 );
557 assert_eq!(stats.min_i64.unwrap(), base);
558 assert_eq!(stats.max_i64.unwrap(), base + 9 * 100_000);
559
560 assert!(
562 !ScanPredicate::eq_i64(0, target).can_skip_block(stats),
563 "must not skip: target={target} is within the block range"
564 );
565
566 assert!(
568 ScanPredicate::eq_i64(0, base - 1).can_skip_block(stats),
569 "must skip: base-1 is below block min"
570 );
571
572 let min_f64 = base as f64;
575 let target_f64 = target as f64;
576 let max_f64 = (base + 9 * 100_000) as f64;
577 let _ = (min_f64, target_f64, max_f64); }
581
582 #[test]
583 fn integer_block_stats_have_exact_i64_fields() {
584 let schema = ColumnarSchema::new(vec![
586 ColumnDef::required("id", ColumnType::Int64).with_primary_key(),
587 ])
588 .expect("valid");
589
590 let mut mt = ColumnarMemtable::new(&schema);
591 for i in 0..5i64 {
592 mt.append_row(&[Value::Integer(i64::MAX - 4 + i)])
593 .expect("append");
594 }
595
596 let (schema, columns, row_count) = mt.drain();
597 let segment = SegmentWriter::plain()
598 .write_segment(&schema, &columns, row_count, None)
599 .expect("write");
600 let footer = SegmentFooter::from_segment_tail(&segment).expect("footer");
601
602 let stats = &footer.columns[0].block_stats[0];
603 assert_eq!(stats.min_i64, Some(i64::MAX - 4));
604 assert_eq!(stats.max_i64, Some(i64::MAX));
605
606 use crate::predicate::ScanPredicate;
608 assert!(!ScanPredicate::eq_i64(0, i64::MAX - 2).can_skip_block(stats));
609 assert!(ScanPredicate::eq_i64(0, i64::MAX - 10).can_skip_block(stats));
611 }
612
613 #[test]
614 fn string_block_stats_bloom_rejects_absent_value() {
615 let schema = ColumnarSchema::new(vec![ColumnDef::required("label", ColumnType::String)])
616 .expect("valid");
617
618 let mut mt = ColumnarMemtable::new(&schema);
619 let values: Vec<String> = (0..20).map(|i| format!("val_{i:02}")).collect();
621 for name in &values {
622 mt.append_row(&[Value::String(name.clone())])
623 .expect("append");
624 }
625 mt.append_row(&[Value::String("alpha".into())])
627 .expect("append");
628 mt.append_row(&[Value::String("beta".into())])
629 .expect("append");
630 mt.append_row(&[Value::String("gamma".into())])
631 .expect("append");
632
633 let (schema, columns, row_count) = mt.drain();
634 let segment = SegmentWriter::plain()
635 .write_segment(&schema, &columns, row_count, None)
636 .expect("write");
637 let footer = SegmentFooter::from_segment_tail(&segment).expect("footer");
638 let stats = &footer.columns[0].block_stats[0];
639
640 use crate::predicate::{ScanPredicate, bloom_may_contain};
641
642 let bloom = stats
643 .bloom
644 .as_ref()
645 .expect("bloom present for >16 distinct");
646 assert!(bloom_may_contain(bloom, "alpha"));
647 assert!(bloom_may_contain(bloom, "beta"));
648 assert!(bloom_may_contain(bloom, "gamma"));
649
650 let delta_absent = !bloom_may_contain(bloom, "delta");
652 if delta_absent {
653 assert!(ScanPredicate::str_eq(0, "delta").can_skip_block(stats));
655 }
656 }
657}