Skip to main content

reifydb_codec/frame/decode/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4//! Decoder side of RBCF, reversing `encode/` header by header. Payloads can arrive from untrusted
5//! peers, so every malformed input returns a typed `DecodeError` rather than panicking.
6
7mod 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}