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