1mod any;
8mod fixed;
9mod varlen;
10
11use std::str;
12
13use bigdecimal::BigDecimal;
14use num_bigint::BigInt;
15use reifydb_value::{
16 encoding::LeBytes,
17 reifydb_assertions,
18 util::bitvec::BitVec,
19 value::{
20 container::{blob::BlobContainer, number::NumberContainer, utf8::Utf8Container},
21 datetime::DateTime,
22 decimal::Decimal,
23 frame::{column::FrameColumn, data::FrameColumnData, frame::Frame},
24 int::Int,
25 row_number::RowNumber,
26 system_columns::SystemColumns,
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_TIME, META_HAS_UPDATED_AT, RBCF_MAGIC,
39 RBCF_VERSION, 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 (time, pos) = read_datetime_array(data, pos, header.row_count, header.meta_flags, META_HAS_TIME)?;
99 let (columns, pos) = read_frame_columns(data, pos, header.column_count)?;
100
101 Ok((
102 Frame {
103 system: SystemColumns::new(row_numbers, Vec::new(), created_at, updated_at, time),
104 columns,
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 _reserved = 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 },
130 pos,
131 ))
132}
133
134#[inline]
135fn read_row_numbers(
136 data: &[u8],
137 mut pos: usize,
138 row_count: usize,
139 meta_flags: u8,
140) -> Result<(Vec<RowNumber>, usize), DecodeError> {
141 if meta_flags & META_HAS_ROW_NUMBERS == 0 {
142 return Ok((Vec::new(), pos));
143 }
144 check_len(data, pos, row_count * RowNumber::ENCODED_SIZE)?;
145 let mut row_numbers = Vec::with_capacity(row_count);
146 for _ in 0..row_count {
147 row_numbers.push(RowNumber::read_le(&data[pos..]));
148 pos += RowNumber::ENCODED_SIZE;
149 }
150 Ok((row_numbers, pos))
151}
152
153#[inline]
154fn read_datetime_array(
155 data: &[u8],
156 mut pos: usize,
157 row_count: usize,
158 meta_flags: u8,
159 flag: u8,
160) -> Result<(Vec<DateTime>, usize), DecodeError> {
161 if meta_flags & flag == 0 {
162 return Ok((Vec::new(), pos));
163 }
164 check_len(data, pos, row_count * DateTime::ENCODED_SIZE)?;
165 let mut values = Vec::with_capacity(row_count);
166 for _ in 0..row_count {
167 values.push(DateTime::read_le(&data[pos..]));
168 pos += DateTime::ENCODED_SIZE;
169 }
170 Ok((values, pos))
171}
172
173#[inline]
174fn read_frame_columns(
175 data: &[u8],
176 mut pos: usize,
177 column_count: usize,
178) -> Result<(Vec<FrameColumn>, usize), DecodeError> {
179 let mut columns = Vec::with_capacity(column_count);
180 for _ in 0..column_count {
181 let (col, new_pos) = decode_column(data, pos)?;
182 columns.push(col);
183 pos = new_pos;
184 }
185 Ok((columns, pos))
186}
187
188fn decode_column(data: &[u8], start: usize) -> Result<(FrameColumn, usize), DecodeError> {
189 let mut pos = start;
190 check_len(data, pos, COLUMN_DESCRIPTOR_SIZE)?;
191
192 let type_code = data[pos];
193 pos += 1;
194 let encoding_byte = data[pos];
195 pos += 1;
196 let flags = data[pos];
197 pos += 1;
198 let _reserved = data[pos];
199 pos += 1;
200 let name_len = read_u16(data, pos) as usize;
201 pos += 2;
202 let _reserved2 = read_u16(data, pos);
203 pos += 2;
204 let row_count = read_u32(data, pos) as usize;
205 pos += 4;
206 let nones_len = read_u32(data, pos) as usize;
207 pos += 4;
208 let data_len = read_u32(data, pos) as usize;
209 pos += 4;
210 let offsets_len = read_u32(data, pos) as usize;
211 pos += 4;
212 let extra_len = read_u32(data, pos) as usize;
213 pos += 4;
214
215 let encoding = Encoding::from_u8(encoding_byte).ok_or(DecodeError::UnknownEncoding(encoding_byte))?;
216 let has_nones = flags & COL_FLAG_HAS_NONES != 0;
217
218 check_len(data, pos, name_len)?;
219 let name = str::from_utf8(&data[pos..pos + name_len])
220 .map_err(|e| DecodeError::InvalidData(format!("invalid column name: {}", e)))?
221 .to_string();
222 pos += name_len;
223 let name_pad = (4 - (name_len % 4)) % 4;
224 pos += name_pad;
225
226 let result = (|| -> Result<(FrameColumnData, usize), DecodeError> {
227 let mut pos = pos;
228
229 let nones = if has_nones && nones_len > 0 {
230 check_len(data, pos, nones_len)?;
231 let bv = decode_bitvec(&data[pos..pos + nones_len], row_count);
232 pos += nones_len;
233 Some(bv)
234 } else {
235 pos += nones_len;
236 None
237 };
238
239 check_len(data, pos, data_len)?;
240 let data_bytes = &data[pos..pos + data_len];
241 pos += data_len;
242
243 check_len(data, pos, offsets_len)?;
244 let offsets_bytes = &data[pos..pos + offsets_len];
245 pos += offsets_len;
246
247 check_len(data, pos, extra_len)?;
248 let extra_bytes = &data[pos..pos + extra_len];
249 pos += extra_len;
250
251 let col_data = decode_column_dispatch(
252 type_code,
253 encoding,
254 flags,
255 row_count,
256 data_bytes,
257 offsets_bytes,
258 extra_bytes,
259 )?;
260
261 let col_data = if let Some(bitvec) = nones {
262 FrameColumnData::Option {
263 inner: Box::new(col_data),
264 bitvec,
265 }
266 } else {
267 col_data
268 };
269
270 Ok((col_data, pos))
271 })()
272 .map_err(|e| DecodeError::ColumnDecodeFailed {
273 column_name: name.clone(),
274 row_index: None,
275 source: Box::new(e),
276 })?;
277
278 let (col_data, pos) = result;
279 Ok((
280 FrameColumn {
281 name,
282 data: col_data,
283 },
284 pos,
285 ))
286}
287
288pub(crate) fn column_type_from_code(type_code: u8) -> Result<ValueType, DecodeError> {
289 let tag = TypeTag::from_byte(type_code)?;
290 if tag.depth() != 0 {
291 return Err(DecodeError::InvalidData(format!(
292 "column type code 0x{type_code:02X} carries option depth"
293 )));
294 }
295 tag.to_type()
296}
297
298fn decode_column_dispatch(
299 type_code: u8,
300 encoding: Encoding,
301 flags: u8,
302 row_count: usize,
303 data: &[u8],
304 offsets: &[u8],
305 extra: &[u8],
306) -> Result<FrameColumnData, DecodeError> {
307 let ty = column_type_from_code(type_code)?;
308
309 match encoding {
310 Encoding::Plain | Encoding::BitPack => {
311 if ty == ValueType::Any {
312 return any::decode_any_column(row_count, data);
313 }
314
315 if let Some(result) = fixed::decode_fixed_plain(type_code, row_count, data) {
316 return result;
317 }
318
319 if let Some(result) = varlen::decode_varlen_plain(type_code, row_count, data, offsets) {
320 return result;
321 }
322 Err(DecodeError::UnsupportedType(format!("{:?}", ty)))
323 }
324 Encoding::Dict => match ty {
325 ValueType::Utf8 => {
326 let index_width = dict_index_width_from_flags(flags);
327 let strings = decode_dict_utf8(data, extra, row_count, index_width)?;
328 Ok(FrameColumnData::Utf8(Utf8Container::new(strings)))
329 }
330 ValueType::Blob => {
331 let index_width = dict_index_width_from_flags(flags);
332 let blobs = decode_dict_blob(data, extra, row_count, index_width)?;
333 Ok(FrameColumnData::Blob(BlobContainer::new(blobs)))
334 }
335 ValueType::Int => {
336 let index_width = dict_index_width_from_flags(flags);
337 let dict_entries = decode_dict_table_bytes(extra)?;
338 let mut values = Vec::with_capacity(row_count);
339 for i in 0..row_count {
340 let idx = read_index(data, i, index_width) as usize;
341 if idx >= dict_entries.len() {
342 return Err(DecodeError::InvalidData(format!(
343 "dict index {} out of range (dict has {} entries)",
344 idx,
345 dict_entries.len()
346 )));
347 }
348 let big = BigInt::from_signed_bytes_le(&dict_entries[idx]);
349 values.push(Int(big));
350 }
351 Ok(FrameColumnData::Int(NumberContainer::new(values)))
352 }
353 ValueType::Uint => {
354 let index_width = dict_index_width_from_flags(flags);
355 let dict_entries = decode_dict_table_bytes(extra)?;
356 let mut values = Vec::with_capacity(row_count);
357 for i in 0..row_count {
358 let idx = read_index(data, i, index_width) as usize;
359 if idx >= dict_entries.len() {
360 return Err(DecodeError::InvalidData(format!(
361 "dict index {} out of range (dict has {} entries)",
362 idx,
363 dict_entries.len()
364 )));
365 }
366 let big = BigInt::from_signed_bytes_le(&dict_entries[idx]);
367 values.push(Uint(big));
368 }
369 Ok(FrameColumnData::Uint(NumberContainer::new(values)))
370 }
371 ValueType::Decimal => {
372 let index_width = dict_index_width_from_flags(flags);
373 let dict_entries = decode_dict_table_bytes(extra)?;
374 let mut values = Vec::with_capacity(row_count);
375 for i in 0..row_count {
376 let idx = read_index(data, i, index_width) as usize;
377 if idx >= dict_entries.len() {
378 return Err(DecodeError::InvalidData(format!(
379 "dict index {} out of range (dict has {} entries)",
380 idx,
381 dict_entries.len()
382 )));
383 }
384 let s = str::from_utf8(&dict_entries[idx]).map_err(|e| {
385 DecodeError::InvalidData(format!("invalid decimal string: {}", e))
386 })?;
387 let dec: BigDecimal = s.parse().map_err(|e| {
388 DecodeError::InvalidData(format!("invalid decimal: {}", e))
389 })?;
390 values.push(Decimal::new(dec));
391 }
392 Ok(FrameColumnData::Decimal(NumberContainer::new(values)))
393 }
394 _ => Err(DecodeError::InvalidData(format!("Dict encoding not supported for type {:?}", ty))),
395 },
396 Encoding::Rle => match ty {
397 ValueType::Int | ValueType::Uint | ValueType::Decimal => {
398 varlen::decode_rle_varlen_column(type_code, row_count, data)
399 }
400 _ => fixed::decode_rle_column(type_code, row_count, data),
401 },
402 Encoding::Delta => fixed::decode_delta_column(type_code, row_count, data),
403 Encoding::DeltaRle => fixed::decode_delta_rle_column(type_code, row_count, data),
404 }
405}
406
407fn decode_bitvec(data: &[u8], len: usize) -> BitVec {
408 BitVec::from_raw(data.to_vec(), len)
409}
410
411#[inline]
412fn read_u16(data: &[u8], pos: usize) -> u16 {
413 u16::read_le(&data[pos..])
414}
415
416#[inline]
417fn read_u32(data: &[u8], pos: usize) -> u32 {
418 u32::read_le(&data[pos..])
419}
420
421fn check_len(data: &[u8], pos: usize, needed: usize) -> Result<(), DecodeError> {
422 if pos + needed > data.len() {
423 Err(DecodeError::UnexpectedEof {
424 expected: needed,
425 available: data.len().saturating_sub(pos),
426 })
427 } else {
428 Ok(())
429 }
430}