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
//! Compression codecs for native protocol frames.
//!
//! Built-in LZ4 and Zstandard codecs require corresponding crate features.
//! [`Codec::empty`] and [`Codec::from_raw`] support custom implementations.
use core::pin::Pin;
use crate::sys;
/// Compression algorithm used for native protocol frames.
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
#[repr(i32)]
pub enum Compression {
/// Disables compression and does not require a [`Codec`].
#[default]
None = sys::CHC_COMP_NONE,
/// Uses LZ4 and requires corresponding callbacks, such as [`Codec::lz4`].
Lz4 = sys::CHC_COMP_LZ4,
/// Uses Zstandard and requires corresponding callbacks, such as
/// [`Codec::zstd`].
Zstd = sys::CHC_COMP_ZSTD,
}
/// Compression callbacks used by clickhouse-c.
///
/// Value is pinned because C code retains address of callback table.
pub struct Codec {
raw: sys::chc_codec,
_pin: core::marker::PhantomPinned,
}
impl Codec {
/// Creates a codec without compression callbacks.
///
/// Codec only supports [`Compression::None`] until required callbacks are
/// installed through [`raw_mut`]. [`Client::init`](crate::Client::init)
/// rejects missing callbacks.
///
/// [`raw_mut`]: Codec::raw_mut
pub fn empty() -> Pin<Box<Self>> {
Box::pin(Self {
raw: sys::chc_codec {
ud: core::ptr::null_mut(),
lz4_compress: None,
lz4_decompress: None,
zstd_compress: None,
zstd_decompress: None,
lz4_bound: None,
zstd_bound: None,
},
_pin: core::marker::PhantomPinned,
})
}
/// Creates a codec from a raw callback table.
///
/// # Safety
///
/// Each function pointer must match corresponding field signature. All
/// callbacks required by selected [`Compression`] must be present. User
/// data referenced by `ud` must remain valid while codec exists and from
/// every thread that uses codec.
pub unsafe fn from_raw(raw: sys::chc_codec) -> Pin<Box<Self>> {
Box::pin(Self {
raw,
_pin: core::marker::PhantomPinned,
})
}
/// Creates built-in LZ4 codec backed by system liblz4.
#[cfg(feature = "lz4")]
pub fn lz4() -> Pin<Box<Self>> {
let mut b = Self::empty();
unsafe {
let this = b.as_mut().get_unchecked_mut();
sys::chc_lz4_codec_init(&mut this.raw);
}
b
}
/// Creates built-in Zstandard codec backed by system libzstd.
#[cfg(feature = "zstd")]
pub fn zstd() -> Pin<Box<Self>> {
let mut b = Self::empty();
unsafe {
let this = b.as_mut().get_unchecked_mut();
sys::chc_zstd_codec_init(&mut this.raw);
}
b
}
/// Returns mutable access to raw callback table.
///
/// # Safety
///
/// C library calls installed function pointers without validation.
///
/// * Each function pointer must match corresponding field signature.
/// * All callbacks required by selected [`Compression`] must be present.
/// For example, LZ4 requires `lz4_compress`, `lz4_decompress`, and
/// `lz4_bound`.
/// * Data referenced by `ud` must remain valid while codec exists and from
/// every thread that uses codec.
pub unsafe fn raw_mut(self: Pin<&mut Self>) -> &mut sys::chc_codec {
unsafe { &mut self.get_unchecked_mut().raw }
}
#[inline]
pub(crate) fn as_ptr(self: Pin<&Self>) -> *const sys::chc_codec {
&self.raw
}
/// Returns whether all callbacks required by `compression` are present.
pub(crate) fn supports(self: Pin<&Self>, compression: Compression) -> bool {
match compression {
Compression::None => true,
Compression::Lz4 => {
self.raw.lz4_compress.is_some()
&& self.raw.lz4_decompress.is_some()
&& self.raw.lz4_bound.is_some()
}
Compression::Zstd => {
self.raw.zstd_compress.is_some()
&& self.raw.zstd_decompress.is_some()
&& self.raw.zstd_bound.is_some()
}
}
}
}
unsafe impl Send for Codec {}
/// Calculates CityHash128 and returns low and high words in wire order.
pub fn cityhash128(data: &[u8]) -> (u64, u64) {
let mut lo = 0u64;
let mut hi = 0u64;
unsafe {
sys::chc_cityhash128(data.as_ptr().cast(), data.len(), &mut lo, &mut hi);
}
(lo, hi)
}
#[cfg(test)]
mod tests {
use core::ffi::{c_int, c_void};
use super::{Codec, Compression, cityhash128};
use crate::sys;
#[test]
fn empty_codec_supports_only_uncompressed() {
let codec = Codec::empty();
assert!(codec.as_ref().supports(Compression::None));
assert!(!codec.as_ref().supports(Compression::Lz4));
assert!(!codec.as_ref().supports(Compression::Zstd));
}
// Frame allocation requires bound callback before compression
#[test]
fn missing_bound_callback_is_not_support() {
let mut codec = Codec::empty();
unsafe {
let raw = codec.as_mut().raw_mut();
raw.lz4_compress = Some(stub_compress);
raw.lz4_decompress = Some(stub_decompress);
}
assert!(!codec.as_ref().supports(Compression::Lz4));
unsafe { codec.as_mut().raw_mut().lz4_bound = Some(stub_bound) };
assert!(codec.as_ref().supports(Compression::Lz4));
// Installed pointers are what C would call
let raw = unsafe { codec.as_mut().raw_mut() };
let mut err = sys::chc_err::zeroed();
let mut written = 0usize;
assert_eq!(unsafe { raw.lz4_bound.expect("bound")(7) }, 7);
assert_eq!(
unsafe {
raw.lz4_compress.expect("compress")(
raw.ud,
core::ptr::null(),
0,
core::ptr::null_mut(),
0,
&mut written,
&mut err,
)
},
sys::CHC_OK,
);
assert_eq!(
unsafe {
raw.lz4_decompress.expect("decompress")(
raw.ud,
core::ptr::null(),
0,
core::ptr::null_mut(),
0,
&mut err,
)
},
sys::CHC_OK,
);
}
// Custom codecs arrive as a filled table rather than through raw_mut
#[test]
fn from_raw_keeps_every_callback() {
let table = sys::chc_codec {
ud: core::ptr::null_mut(),
lz4_compress: Some(stub_compress),
lz4_decompress: Some(stub_decompress),
lz4_bound: Some(stub_bound),
zstd_compress: Some(stub_compress),
zstd_decompress: Some(stub_decompress),
zstd_bound: Some(stub_bound),
};
let codec = unsafe { Codec::from_raw(table) };
assert!(codec.as_ref().supports(Compression::Lz4));
assert!(codec.as_ref().supports(Compression::Zstd));
assert!(!codec.as_ref().as_ptr().is_null());
}
// Wire order is low word first. Vectors match clickhouse-cpp CityHash128,
// which is what a server checksums against
#[test]
fn cityhash128_matches_the_reference_digest() {
assert_eq!(cityhash128(b""), (0x3df09dfc64c09a2b, 0x3cb540c392e51e29));
assert_eq!(
cityhash128(b"Hello, World!"),
(0x703dabf8d081ec00, 0xa196e28f28c3ee09),
);
// 200 bytes reach the unrolled path for inputs of 128 bytes and more
assert_eq!(
cityhash128(&[0xab; 200]),
(0xb29e1d196fe650df, 0x3dae10d6a77e0432),
);
}
// Test only checks whether callback is present
unsafe extern "C" fn stub_compress(
_ud: *mut c_void,
_src: *const c_void,
_src_len: usize,
_dst: *mut c_void,
_dst_cap: usize,
_dst_n: *mut usize,
_err: *mut sys::chc_err,
) -> c_int {
sys::CHC_OK
}
unsafe extern "C" fn stub_decompress(
_ud: *mut c_void,
_src: *const c_void,
_src_len: usize,
_dst: *mut c_void,
_original_size: usize,
_err: *mut sys::chc_err,
) -> c_int {
sys::CHC_OK
}
unsafe extern "C" fn stub_bound(src_len: usize) -> usize {
src_len
}
#[cfg(feature = "lz4")]
#[test]
fn built_in_lz4_fills_every_slot() {
assert!(Codec::lz4().as_ref().supports(Compression::Lz4));
}
#[cfg(feature = "zstd")]
#[test]
fn built_in_zstd_fills_every_slot() {
assert!(Codec::zstd().as_ref().supports(Compression::Zstd));
}
}