1mod any;
5mod fixed;
6mod varlen;
7
8use std::str;
9
10use bigdecimal::BigDecimal;
11use num_bigint::BigInt;
12use reifydb_value::{
13 encoding::LeBytes,
14 reifydb_assertions,
15 util::bitvec::BitVec,
16 value::{
17 container::{blob::BlobContainer, number::NumberContainer, utf8::Utf8Container},
18 datetime::DateTime,
19 decimal::Decimal,
20 diff_type::DiffType,
21 frame::{column::FrameColumn, data::FrameColumnData, frame::Frame},
22 int::Int,
23 row_number::RowNumber,
24 system_columns::SystemColumns,
25 uint::Uint,
26 value_type::ValueType,
27 },
28};
29
30use crate::{
31 error::DecodeError,
32 frame::{
33 encoding::dict::{decode_dict_blob, decode_dict_table_bytes, decode_dict_utf8, read_index},
34 format::{
35 COL_FLAG_HAS_NONES, COLUMN_DESCRIPTOR_SIZE, Encoding, FRAME_HEADER_SIZE, MESSAGE_HEADER_SIZE,
36 META_HAS_CREATED_AT, META_HAS_ROW_NUMBERS, META_HAS_TIME, META_HAS_UPDATED_AT, RBCF_MAGIC,
37 RBCF_VERSION, dict_index_width_from_flags,
38 },
39 },
40 tag::TypeTag,
41};
42
43pub fn decode_frames(data: &[u8]) -> Result<Vec<Frame>, DecodeError> {
44 let mut pos = 0;
45
46 check_len(data, pos, MESSAGE_HEADER_SIZE)?;
47 let magic = read_u32(data, pos);
48 pos += 4;
49 if magic != RBCF_MAGIC {
50 return Err(DecodeError::InvalidMagic(magic));
51 }
52 let version = read_u16(data, pos);
53 pos += 2;
54 if version != RBCF_VERSION {
55 return Err(DecodeError::UnsupportedVersion(version));
56 }
57 let _flags = read_u16(data, pos);
58 pos += 2;
59 let frame_count = read_u32(data, pos) as usize;
60 pos += 4;
61 let _total_size = read_u32(data, pos) as usize;
62 pos += 4;
63
64 let mut frames = Vec::with_capacity(frame_count);
65 for _ in 0..frame_count {
66 let (frame, new_pos) = decode_frame(data, pos)?;
67 frames.push(frame);
68 pos = new_pos;
69 }
70
71 reifydb_assertions! {
72 assert!(
73 pos == _total_size,
74 "the RBCF message header declared a total size that disagrees with the bytes consumed while decoding, so the message is truncated or carries trailing bytes and a peer would mis-frame the following message (declared={} consumed={})",
75 _total_size,
76 pos
77 );
78 }
79
80 Ok(frames)
81}
82
83struct FrameHeader {
84 row_count: usize,
85 column_count: usize,
86 meta_flags: u8,
87 op: Option<DiffType>,
88}
89
90fn decode_frame(data: &[u8], start: usize) -> Result<(Frame, usize), DecodeError> {
91 let (header, pos) = read_frame_header(data, start)?;
92 let (row_numbers, pos) = read_row_numbers(data, pos, header.row_count, header.meta_flags)?;
93 let (created_at, pos) =
94 read_datetime_array(data, pos, header.row_count, header.meta_flags, META_HAS_CREATED_AT)?;
95 let (updated_at, pos) =
96 read_datetime_array(data, pos, header.row_count, header.meta_flags, META_HAS_UPDATED_AT)?;
97 let (time, pos) = read_datetime_array(data, pos, header.row_count, header.meta_flags, META_HAS_TIME)?;
98 let (columns, pos) = read_frame_columns(data, pos, header.column_count)?;
99
100 Ok((
101 Frame {
102 system: SystemColumns::new(row_numbers, Vec::new(), created_at, updated_at, time),
103 columns,
104 op: header.op,
105 },
106 pos,
107 ))
108}
109
110#[inline]
111fn read_frame_header(data: &[u8], start: usize) -> Result<(FrameHeader, usize), DecodeError> {
112 let mut pos = start;
113 check_len(data, pos, FRAME_HEADER_SIZE)?;
114 let row_count = read_u32(data, pos) as usize;
115 pos += 4;
116 let column_count = read_u16(data, pos) as usize;
117 pos += 2;
118 let meta_flags = data[pos];
119 pos += 1;
120 let op = DiffType::from_u8(data[pos]);
121 pos += 1;
122 let _frame_size = read_u32(data, pos);
123 pos += 4;
124 Ok((
125 FrameHeader {
126 row_count,
127 column_count,
128 meta_flags,
129 op,
130 },
131 pos,
132 ))
133}
134
135#[inline]
136fn read_row_numbers(
137 data: &[u8],
138 mut pos: usize,
139 row_count: usize,
140 meta_flags: u8,
141) -> Result<(Vec<RowNumber>, usize), DecodeError> {
142 if meta_flags & META_HAS_ROW_NUMBERS == 0 {
143 return Ok((Vec::new(), pos));
144 }
145 check_len(data, pos, row_count * RowNumber::ENCODED_SIZE)?;
146 let mut row_numbers = Vec::with_capacity(row_count);
147 for _ in 0..row_count {
148 row_numbers.push(RowNumber::read_le(&data[pos..]));
149 pos += RowNumber::ENCODED_SIZE;
150 }
151 Ok((row_numbers, pos))
152}
153
154#[inline]
155fn read_datetime_array(
156 data: &[u8],
157 mut pos: usize,
158 row_count: usize,
159 meta_flags: u8,
160 flag: u8,
161) -> Result<(Vec<DateTime>, usize), DecodeError> {
162 if meta_flags & flag == 0 {
163 return Ok((Vec::new(), pos));
164 }
165 check_len(data, pos, row_count * DateTime::ENCODED_SIZE)?;
166 let mut values = Vec::with_capacity(row_count);
167 for _ in 0..row_count {
168 values.push(DateTime::read_le(&data[pos..]));
169 pos += DateTime::ENCODED_SIZE;
170 }
171 Ok((values, pos))
172}
173
174#[inline]
175fn read_frame_columns(
176 data: &[u8],
177 mut pos: usize,
178 column_count: usize,
179) -> Result<(Vec<FrameColumn>, usize), DecodeError> {
180 let mut columns = Vec::with_capacity(column_count);
181 for _ in 0..column_count {
182 let (col, new_pos) = decode_column(data, pos)?;
183 columns.push(col);
184 pos = new_pos;
185 }
186 Ok((columns, pos))
187}
188
189fn decode_column(data: &[u8], start: usize) -> Result<(FrameColumn, usize), DecodeError> {
190 let mut pos = start;
191 check_len(data, pos, COLUMN_DESCRIPTOR_SIZE)?;
192
193 let type_code = data[pos];
194 pos += 1;
195 let encoding_byte = data[pos];
196 pos += 1;
197 let flags = data[pos];
198 pos += 1;
199 let _reserved = data[pos];
200 pos += 1;
201 let name_len = read_u16(data, pos) as usize;
202 pos += 2;
203 let _reserved2 = read_u16(data, pos);
204 pos += 2;
205 let row_count = read_u32(data, pos) as usize;
206 pos += 4;
207 let nones_len = read_u32(data, pos) as usize;
208 pos += 4;
209 let data_len = read_u32(data, pos) as usize;
210 pos += 4;
211 let offsets_len = read_u32(data, pos) as usize;
212 pos += 4;
213 let extra_len = read_u32(data, pos) as usize;
214 pos += 4;
215
216 let encoding = Encoding::from_u8(encoding_byte).ok_or(DecodeError::UnknownEncoding(encoding_byte))?;
217 let has_nones = flags & COL_FLAG_HAS_NONES != 0;
218
219 check_len(data, pos, name_len)?;
220 let name = str::from_utf8(&data[pos..pos + name_len])
221 .map_err(|e| DecodeError::InvalidData(format!("invalid column name: {}", e)))?
222 .to_string();
223 pos += name_len;
224 let name_pad = (4 - (name_len % 4)) % 4;
225 pos += name_pad;
226
227 let result = (|| -> Result<(FrameColumnData, usize), DecodeError> {
228 let mut pos = pos;
229
230 let nones = if has_nones && nones_len > 0 {
231 check_len(data, pos, nones_len)?;
232 let bv = decode_bitvec(&data[pos..pos + nones_len], row_count);
233 pos += nones_len;
234 Some(bv)
235 } else {
236 pos += nones_len;
237 None
238 };
239
240 check_len(data, pos, data_len)?;
241 let data_bytes = &data[pos..pos + data_len];
242 pos += data_len;
243
244 check_len(data, pos, offsets_len)?;
245 let offsets_bytes = &data[pos..pos + offsets_len];
246 pos += offsets_len;
247
248 check_len(data, pos, extra_len)?;
249 let extra_bytes = &data[pos..pos + extra_len];
250 pos += extra_len;
251
252 let col_data = decode_column_dispatch(
253 type_code,
254 encoding,
255 flags,
256 row_count,
257 data_bytes,
258 offsets_bytes,
259 extra_bytes,
260 )?;
261
262 let col_data = if let Some(bitvec) = nones {
263 FrameColumnData::Option {
264 inner: Box::new(col_data),
265 bitvec,
266 }
267 } else {
268 col_data
269 };
270
271 Ok((col_data, pos))
272 })()
273 .map_err(|e| DecodeError::ColumnDecodeFailed {
274 column_name: name.clone(),
275 row_index: None,
276 source: Box::new(e),
277 })?;
278
279 let (col_data, pos) = result;
280 Ok((
281 FrameColumn {
282 name,
283 data: col_data,
284 },
285 pos,
286 ))
287}
288
289pub(crate) fn column_type_from_code(type_code: u8) -> Result<ValueType, DecodeError> {
290 let tag = TypeTag::from_byte(type_code)?;
291 if tag.depth() != 0 {
292 return Err(DecodeError::InvalidData(format!(
293 "column type code 0x{type_code:02X} carries option depth"
294 )));
295 }
296 tag.to_type()
297}
298
299fn decode_column_dispatch(
300 type_code: u8,
301 encoding: Encoding,
302 flags: u8,
303 row_count: usize,
304 data: &[u8],
305 offsets: &[u8],
306 extra: &[u8],
307) -> Result<FrameColumnData, DecodeError> {
308 let ty = column_type_from_code(type_code)?;
309
310 match encoding {
311 Encoding::Plain | Encoding::BitPack => {
312 if ty == ValueType::Any {
313 return any::decode_any_column(row_count, data);
314 }
315
316 if let Some(result) = fixed::decode_fixed_plain(type_code, row_count, data) {
317 return result;
318 }
319
320 if let Some(result) = varlen::decode_varlen_plain(type_code, row_count, data, offsets) {
321 return result;
322 }
323 Err(DecodeError::UnsupportedType(format!("{:?}", ty)))
324 }
325 Encoding::Dict => match ty {
326 ValueType::Utf8 => {
327 let index_width = dict_index_width_from_flags(flags);
328 let strings = decode_dict_utf8(data, extra, row_count, index_width)?;
329 Ok(FrameColumnData::Utf8(Utf8Container::new(strings)))
330 }
331 ValueType::Blob => {
332 let index_width = dict_index_width_from_flags(flags);
333 let blobs = decode_dict_blob(data, extra, row_count, index_width)?;
334 Ok(FrameColumnData::Blob(BlobContainer::new(blobs)))
335 }
336 ValueType::Int => {
337 let index_width = dict_index_width_from_flags(flags);
338 let dict_entries = decode_dict_table_bytes(extra)?;
339 let mut values = Vec::with_capacity(row_count);
340 for i in 0..row_count {
341 let idx = read_index(data, i, index_width) as usize;
342 if idx >= dict_entries.len() {
343 return Err(DecodeError::InvalidData(format!(
344 "dict index {} out of range (dict has {} entries)",
345 idx,
346 dict_entries.len()
347 )));
348 }
349 let big = BigInt::from_signed_bytes_le(&dict_entries[idx]);
350 values.push(Int(big));
351 }
352 Ok(FrameColumnData::Int(NumberContainer::new(values)))
353 }
354 ValueType::Uint => {
355 let index_width = dict_index_width_from_flags(flags);
356 let dict_entries = decode_dict_table_bytes(extra)?;
357 let mut values = Vec::with_capacity(row_count);
358 for i in 0..row_count {
359 let idx = read_index(data, i, index_width) as usize;
360 if idx >= dict_entries.len() {
361 return Err(DecodeError::InvalidData(format!(
362 "dict index {} out of range (dict has {} entries)",
363 idx,
364 dict_entries.len()
365 )));
366 }
367 let big = BigInt::from_signed_bytes_le(&dict_entries[idx]);
368 values.push(Uint(big));
369 }
370 Ok(FrameColumnData::Uint(NumberContainer::new(values)))
371 }
372 ValueType::Decimal => {
373 let index_width = dict_index_width_from_flags(flags);
374 let dict_entries = decode_dict_table_bytes(extra)?;
375 let mut values = Vec::with_capacity(row_count);
376 for i in 0..row_count {
377 let idx = read_index(data, i, index_width) as usize;
378 if idx >= dict_entries.len() {
379 return Err(DecodeError::InvalidData(format!(
380 "dict index {} out of range (dict has {} entries)",
381 idx,
382 dict_entries.len()
383 )));
384 }
385 let s = str::from_utf8(&dict_entries[idx]).map_err(|e| {
386 DecodeError::InvalidData(format!("invalid decimal string: {}", e))
387 })?;
388 let dec: BigDecimal = s.parse().map_err(|e| {
389 DecodeError::InvalidData(format!("invalid decimal: {}", e))
390 })?;
391 values.push(Decimal::new(dec));
392 }
393 Ok(FrameColumnData::Decimal(NumberContainer::new(values)))
394 }
395 _ => Err(DecodeError::InvalidData(format!("Dict encoding not supported for type {:?}", ty))),
396 },
397 Encoding::Rle => match ty {
398 ValueType::Int | ValueType::Uint | ValueType::Decimal => {
399 varlen::decode_rle_varlen_column(type_code, row_count, data)
400 }
401 _ => fixed::decode_rle_column(type_code, row_count, data),
402 },
403 Encoding::Delta => fixed::decode_delta_column(type_code, row_count, data),
404 Encoding::DeltaRle => fixed::decode_delta_rle_column(type_code, row_count, data),
405 }
406}
407
408fn decode_bitvec(data: &[u8], len: usize) -> BitVec {
409 BitVec::from_raw(data.to_vec(), len)
410}
411
412#[inline]
413fn read_u16(data: &[u8], pos: usize) -> u16 {
414 u16::read_le(&data[pos..])
415}
416
417#[inline]
418fn read_u32(data: &[u8], pos: usize) -> u32 {
419 u32::read_le(&data[pos..])
420}
421
422fn check_len(data: &[u8], pos: usize, needed: usize) -> Result<(), DecodeError> {
423 if pos + needed > data.len() {
424 Err(DecodeError::UnexpectedEof {
425 expected: needed,
426 available: data.len().saturating_sub(pos),
427 })
428 } else {
429 Ok(())
430 }
431}