Skip to main content

clickhouse_c/
block.rs

1//! Native block reading and column access.
2//!
3//! [`BlockReader`] reads block streams from an [`Io`] implementation. TCP
4//! clients also return [`Block`] values for Data packets.
5
6use core::pin::Pin;
7use core::ptr::NonNull;
8use core::slice;
9
10use crate::alloc::Allocator;
11use crate::error::{Result, check};
12use crate::io::Io;
13use crate::sys;
14use crate::types::TypeRef;
15
16/// Storage layout of a decoded column.
17///
18/// Several ClickHouse types can share a layout. Use
19/// [`Block::column_type`] to determine logical type.
20#[derive(Clone, Copy, Debug)]
21#[repr(i32)]
22pub enum ColumnLayout {
23    /// Fixed-width values stored as contiguous little-endian bytes.
24    Fixed = sys::CHC_COL_FIXED,
25    /// Offsets and byte data used by `String` and string-encoded JSON types.
26    String = sys::CHC_COL_STRING,
27    /// Null map and dense inner column used by `Nullable(T)`.
28    Nullable = sys::CHC_COL_NULLABLE,
29    /// Offsets and element column used by `Array(T)`, `Map(K, V)` as
30    /// `Array(Tuple(K, V))`, and `Nested(...)` as `Array(Tuple(fields))`.
31    Array = sys::CHC_COL_ARRAY,
32    /// Parallel child columns used by tuples, nested values, geographic
33    /// types, and QBit values.
34    Tuple = sys::CHC_COL_TUPLE,
35    /// Keys and dictionary column used by `LowCardinality(T)`.
36    LowCardinality = sys::CHC_COL_LOW_CARDINALITY,
37    /// Empty storage used by `Nothing` and `Nullable(Nothing)` inner columns.
38    Nothing = sys::CHC_COL_NOTHING,
39}
40
41impl ColumnLayout {
42    pub(crate) fn from_raw(k: sys::chc_col_kind) -> Option<Self> {
43        Some(match k {
44            sys::CHC_COL_FIXED => Self::Fixed,
45            sys::CHC_COL_STRING => Self::String,
46            sys::CHC_COL_NULLABLE => Self::Nullable,
47            sys::CHC_COL_ARRAY => Self::Array,
48            sys::CHC_COL_TUPLE => Self::Tuple,
49            sys::CHC_COL_LOW_CARDINALITY => Self::LowCardinality,
50            sys::CHC_COL_NOTHING => Self::Nothing,
51            _ => return None,
52        })
53    }
54}
55
56/// Options that describe Native block framing.
57///
58/// Native files and native TCP protocol use different optional fields. TCP
59/// values depend on negotiated server revision. Incorrect values prevent
60/// block decoding.
61#[derive(Clone, Copy, Default)]
62pub struct BlockOpts {
63    /// Includes 8-byte `BlockInfo` prefix used by TCP revision 51903 and later.
64    pub has_block_info: bool,
65    /// Includes custom serialization flag used by TCP revision 54454 and later.
66    pub has_custom_serialization: bool,
67    /// Read buffer size in bytes. Zero selects 8 KiB default.
68    pub read_buffer_bytes: usize,
69}
70
71impl BlockOpts {
72    pub(crate) fn to_raw(self) -> sys::chc_block_opts {
73        sys::chc_block_opts {
74            has_block_info: self.has_block_info,
75            has_custom_serialization: self.has_custom_serialization,
76            read_buffer_bytes: self.read_buffer_bytes,
77        }
78    }
79}
80
81/// Decoded Native block.
82///
83/// Value releases its memory with allocator used during decoding.
84pub struct Block {
85    raw: NonNull<sys::chc_block>,
86    alloc: Allocator,
87}
88
89impl Block {
90    /// Takes ownership of a raw block pointer.
91    ///
92    /// # Safety
93    /// Caller must own `raw` and stop using it after this call. `alloc` must
94    /// match allocator used to create block.
95    pub(crate) unsafe fn from_raw(raw: *mut sys::chc_block, alloc: Allocator) -> Option<Self> {
96        NonNull::new(raw).map(|raw| Self { raw, alloc })
97    }
98
99    pub fn n_rows(&self) -> usize {
100        unsafe { sys::chc_block_n_rows(self.raw.as_ptr().cast_const()) }
101    }
102
103    pub fn n_columns(&self) -> usize {
104        unsafe { sys::chc_block_n_columns(self.raw.as_ptr().cast_const()) }
105    }
106
107    /// Returns column name bytes without UTF-8 validation.
108    pub fn column_name(&self, i: usize) -> Option<&[u8]> {
109        let mut len = 0;
110        let p = unsafe { sys::chc_block_column_name(self.raw.as_ptr().cast_const(), i, &mut len) };
111        if p.is_null() {
112            None
113        } else {
114            Some(unsafe { slice::from_raw_parts(p.cast::<u8>(), len) })
115        }
116    }
117
118    pub fn column_type(&self, i: usize) -> Option<TypeRef<'_>> {
119        let p = unsafe { sys::chc_block_column_type(self.raw.as_ptr().cast_const(), i) };
120        if p.is_null() {
121            None
122        } else {
123            Some(TypeRef {
124                raw: p,
125                _marker: core::marker::PhantomData,
126            })
127        }
128    }
129
130    pub fn column(&self, i: usize) -> Option<Column<'_>> {
131        let p = unsafe { sys::chc_block_column(self.raw.as_ptr().cast_const(), i) };
132        if p.is_null() {
133            None
134        } else {
135            Some(Column {
136                raw: p,
137                _marker: core::marker::PhantomData,
138            })
139        }
140    }
141
142    /// Validates structural relationships within every column.
143    ///
144    /// Validation checks array offset order and LowCardinality dictionary
145    /// indexes. Call this method before using offsets or keys from untrusted
146    /// data as indexes.
147    ///
148    /// Runtime cost is proportional to row count.
149    pub fn validate(&self) -> Result<()> {
150        for i in 0..self.n_columns() {
151            if let Some(col) = self.column(i) {
152                col.validate()?;
153            }
154        }
155        Ok(())
156    }
157
158    /// Returns server flag marking a block truncated by `max_rows_to_group_by`
159    /// with `group_by_overflow_mode = 'any'`.
160    pub fn is_overflows(&self) -> bool {
161        unsafe { sys::chc_block_is_overflows(self.raw.as_ptr().cast_const()) }
162    }
163
164    /// Returns two-level aggregation bucket, or -1 when not applicable.
165    pub fn bucket_num(&self) -> i32 {
166        unsafe { sys::chc_block_bucket_num(self.raw.as_ptr().cast_const()) }
167    }
168}
169
170impl Drop for Block {
171    fn drop(&mut self) {
172        unsafe { sys::chc_block_destroy(self.raw.as_ptr(), self.alloc.as_ptr()) };
173    }
174}
175
176unsafe impl Send for Block {}
177
178/// Reads consecutive Native [`Block`] values from an [`Io`] implementation.
179///
180/// Reader retains buffered bytes between calls to [`read`](Self::read).
181pub struct BlockReader<'io, I: Io + ?Sized> {
182    raw: NonNull<sys::chc_in>,
183    // `raw` retains pointer into pinned I/O value
184    _io: Pin<&'io mut I>,
185    // C reader retains allocator address until destruction
186    alloc: Box<Allocator>,
187    opts: sys::chc_block_opts,
188}
189
190impl<'io, I: Io + ?Sized> BlockReader<'io, I> {
191    /// Creates a reader using framing described by `opts`.
192    ///
193    /// Use [`BlockOpts::default`] for output from `clickhouse local`.
194    pub fn new(mut io: Pin<&'io mut I>, alloc: Allocator, opts: BlockOpts) -> Result<Self> {
195        let raw_opts = opts.to_raw();
196        // C reader retains this address
197        let alloc = Box::new(alloc);
198        let mut raw: *mut sys::chc_in = core::ptr::null_mut();
199        let mut err = sys::chc_err::zeroed();
200        let rc = unsafe {
201            sys::chc_rs_in_new(
202                io.as_mut().io_ptr(),
203                alloc.as_ptr(),
204                raw_opts.read_buffer_bytes,
205                &mut raw,
206                &mut err,
207            )
208        };
209        check(rc, &err)?;
210        let raw = NonNull::new(raw).expect("chc_rs_in_new returned CHC_OK with null reader");
211        Ok(Self {
212            raw,
213            _io: io,
214            alloc,
215            opts: raw_opts,
216        })
217    }
218
219    /// Decodes next block. Returns `None` for EOF at a block boundary.
220    pub fn read(&mut self) -> Result<Option<Block>> {
221        let mut out: *mut sys::chc_block = core::ptr::null_mut();
222        let mut err = sys::chc_err::zeroed();
223        let rc = unsafe {
224            sys::chc_block_read(
225                self.raw.as_ptr(),
226                self.alloc.as_ptr(),
227                &self.opts,
228                &mut out,
229                &mut err,
230            )
231        };
232        check(rc, &err)?;
233        Ok(NonNull::new(out).map(|raw| Block {
234            raw,
235            alloc: *self.alloc,
236        }))
237    }
238}
239
240impl<I: Io + ?Sized> Drop for BlockReader<'_, I> {
241    fn drop(&mut self) {
242        unsafe { sys::chc_rs_in_destroy(self.raw.as_ptr(), self.alloc.as_ptr()) };
243    }
244}
245
246/// Borrowed view of a block column.
247#[derive(Clone, Copy)]
248pub struct Column<'b> {
249    pub(crate) raw: *const sys::chc_column,
250    pub(crate) _marker: core::marker::PhantomData<&'b sys::chc_column>,
251}
252
253impl<'b> Column<'b> {
254    pub fn layout(&self) -> Option<ColumnLayout> {
255        ColumnLayout::from_raw(unsafe { sys::chc_column_layout(self.raw) })
256    }
257
258    pub fn n_rows(&self) -> usize {
259        unsafe { sys::chc_column_n_rows(self.raw) }
260    }
261
262    /// Validates array offsets and LowCardinality dictionary indexes.
263    ///
264    /// Validation includes nested columns and returns
265    /// [`ErrorKind::Protocol`](crate::ErrorKind::Protocol) for invalid data.
266    pub fn validate(&self) -> Result<()> {
267        let mut err = sys::chc_err::zeroed();
268        let rc = unsafe { sys::chc_column_validate(self.raw, &mut err) };
269        check(rc, &err)
270    }
271
272    /// Returns element width and little-endian data for a fixed-width column.
273    pub fn fixed(&self) -> Option<(usize, &'b [u8])> {
274        let Some(ColumnLayout::Fixed) = self.layout() else {
275            return None;
276        };
277        let mut elem_size = 0usize;
278        let ptr = unsafe { sys::chc_column_fixed_data(self.raw, &mut elem_size) };
279        let n = self.n_rows().checked_mul(elem_size)?;
280        let bytes = unsafe { slice_or_none(ptr.cast::<u8>(), n) }?;
281        Some((elem_size, bytes))
282    }
283
284    /// Returns offsets and bytes for a string column.
285    ///
286    /// Each offset is exclusive end of corresponding row in host byte order.
287    /// Layout also represents string-encoded JSON and LowCardinality
288    /// dictionaries.
289    pub fn string(&self) -> Option<(&'b [u64], &'b [u8])> {
290        let Some(ColumnLayout::String) = self.layout() else {
291            return None;
292        };
293        let n = self.n_rows();
294        let offsets_ptr = unsafe { sys::chc_column_string_offsets(self.raw) };
295        let data_ptr = unsafe { sys::chc_column_string_data(self.raw) };
296        let offsets = unsafe { slice_or_none(offsets_ptr, n) }?;
297        // Bound data by both final offset and recorded allocation size
298        // SAFETY: String layout selects `str_` union member
299        let capacity = unsafe { (*self.raw).payload.str_.bytes };
300        let claimed = offsets.last().copied().unwrap_or(0) as usize;
301        debug_assert!(
302            offsets.windows(2).all(|w| w[0] <= w[1]) && claimed <= capacity,
303            "clickhouse-c published string offsets outside its own data slab",
304        );
305        let data_len = claimed.min(capacity);
306        // Rows of empty strings leave no slab to borrow
307        let data = if data_len == 0 {
308            &[][..]
309        } else {
310            unsafe { slice_or_none(data_ptr, data_len) }?
311        };
312        Some((offsets, data))
313    }
314
315    pub fn null_map(&self) -> Option<&'b [u8]> {
316        let Some(ColumnLayout::Nullable) = self.layout() else {
317            return None;
318        };
319        let p = unsafe { sys::chc_column_null_map(self.raw) };
320        unsafe { slice_or_none(p, self.n_rows()) }
321    }
322
323    pub fn nullable_inner(&self) -> Option<Column<'b>> {
324        let p = unsafe { sys::chc_column_nullable_inner(self.raw) };
325        if p.is_null() {
326            None
327        } else {
328            Some(Column {
329                raw: p,
330                _marker: core::marker::PhantomData,
331            })
332        }
333    }
334
335    pub fn array_offsets(&self) -> Option<&'b [u64]> {
336        let Some(ColumnLayout::Array) = self.layout() else {
337            return None;
338        };
339        let p = unsafe { sys::chc_column_array_offsets(self.raw) };
340        unsafe { slice_or_none(p, self.n_rows()) }
341    }
342
343    pub fn array_values(&self) -> Option<Column<'b>> {
344        let p = unsafe { sys::chc_column_array_values(self.raw) };
345        if p.is_null() {
346            None
347        } else {
348            Some(Column {
349                raw: p,
350                _marker: core::marker::PhantomData,
351            })
352        }
353    }
354
355    pub fn tuple_arity(&self) -> usize {
356        unsafe { sys::chc_column_tuple_arity(self.raw) }
357    }
358
359    pub fn tuple_child(&self, i: usize) -> Option<Column<'b>> {
360        let p = unsafe { sys::chc_column_tuple_child(self.raw, i) };
361        if p.is_null() {
362            None
363        } else {
364            Some(Column {
365                raw: p,
366                _marker: core::marker::PhantomData,
367            })
368        }
369    }
370
371    /// Returns keys and dictionary for a LowCardinality column.
372    pub fn low_cardinality(&self) -> Option<LowCardinalityView<'b>> {
373        let Some(ColumnLayout::LowCardinality) = self.layout() else {
374            return None;
375        };
376        let key_size = unsafe { sys::chc_column_lc_key_size(self.raw) };
377        let key_size = usize::try_from(key_size).ok().filter(|&k| k > 0)?;
378        debug_assert!(
379            matches!(key_size, 1 | 2 | 4 | 8),
380            "clickhouse-c published LowCardinality key_size = {key_size}",
381        );
382        let keys_ptr = unsafe { sys::chc_column_lc_keys(self.raw) };
383        let dict = NonNull::new(unsafe { sys::chc_column_lc_dict(self.raw) }.cast_mut())?;
384        let keys_len = self.n_rows().checked_mul(key_size)?;
385        let keys = unsafe { slice_or_none(keys_ptr.cast::<u8>(), keys_len) }?;
386        Some(LowCardinalityView {
387            key_size,
388            keys,
389            dict: Column {
390                raw: dict.as_ptr(),
391                _marker: core::marker::PhantomData,
392            },
393        })
394    }
395}
396
397/// Borrows `len` elements published by C, or None when C published no slab.
398///
399/// # Safety
400/// `p` must be null or point to `len` initialized elements that outlive `'b`.
401unsafe fn slice_or_none<'b, T>(p: *const T, len: usize) -> Option<&'b [T]> {
402    (!p.is_null()).then(|| unsafe { slice::from_raw_parts(p, len) })
403}
404
405/// Borrowed parts of a LowCardinality column.
406pub struct LowCardinalityView<'b> {
407    /// Key width in bytes. Valid values are 1, 2, 4, and 8.
408    pub key_size: usize,
409    /// Raw keys in host byte order. Length is `n_rows * key_size`.
410    /// Call [`Column::validate`] before using untrusted keys as indexes.
411    pub keys: &'b [u8],
412    /// Dictionary referenced by keys.
413    pub dict: Column<'b>,
414}
415
416#[cfg(test)]
417mod tests {
418    use super::{BlockOpts, ColumnLayout};
419    use crate::sys;
420
421    // Layouts added to C API must not convert to an adjacent Rust variant
422    #[test]
423    fn unknown_layout_is_none() {
424        assert!(ColumnLayout::from_raw(i32::MAX).is_none());
425        assert!(ColumnLayout::from_raw(-1).is_none());
426    }
427
428    #[test]
429    fn every_c_layout_maps_to_its_variant() {
430        for (raw, layout) in [
431            (sys::CHC_COL_FIXED, ColumnLayout::Fixed),
432            (sys::CHC_COL_STRING, ColumnLayout::String),
433            (sys::CHC_COL_NULLABLE, ColumnLayout::Nullable),
434            (sys::CHC_COL_ARRAY, ColumnLayout::Array),
435            (sys::CHC_COL_TUPLE, ColumnLayout::Tuple),
436            (sys::CHC_COL_LOW_CARDINALITY, ColumnLayout::LowCardinality),
437            (sys::CHC_COL_NOTHING, ColumnLayout::Nothing),
438        ] {
439            assert_eq!(
440                ColumnLayout::from_raw(raw).expect("known layout") as i32,
441                layout as i32,
442            );
443        }
444    }
445
446    #[test]
447    fn opts_carry_tcp_framing_into_c() {
448        let raw = BlockOpts {
449            has_block_info: true,
450            has_custom_serialization: true,
451            read_buffer_bytes: 4096,
452        }
453        .to_raw();
454        assert!(raw.has_block_info);
455        assert!(raw.has_custom_serialization);
456        assert_eq!(raw.read_buffer_bytes, 4096);
457    }
458}