Skip to main content

reifydb_codec/frame/encode/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4//! Encoder side of RBCF: frame header, per-column header, optional none-bitmap, encoded value
5//! bytes, with the encoding chosen per column by `heuristics.rs`. The fixed and varlen submodules
6//! cover the two width regimes a column's values can fall into.
7
8pub(crate) mod any;
9mod fixed;
10mod varlen;
11
12use reifydb_value::{
13	encoding::LeBytes,
14	value::frame::{data::FrameColumnData, frame::Frame},
15};
16use tracing::{Span, instrument};
17
18use crate::{
19	error::EncodeError,
20	frame::{
21		encoding::plain::{encode_bitvec, encode_plain},
22		format::{
23			COL_FLAG_HAS_NONES, Encoding, FRAME_HEADER_SIZE, MESSAGE_HEADER_SIZE, META_HAS_CREATED_AT,
24			META_HAS_ROW_NUMBERS, META_HAS_TIME, META_HAS_UPDATED_AT, RBCF_MAGIC, RBCF_VERSION,
25		},
26		heuristics::choose_encoding,
27		options::EncodeOptions,
28	},
29};
30
31pub(crate) struct EncodedColumn {
32	pub(crate) type_code: u8,
33	pub(crate) encoding: Encoding,
34	pub(crate) flags: u8,
35	pub(crate) nones: Vec<u8>,
36	pub(crate) data: Vec<u8>,
37	pub(crate) offsets: Vec<u8>,
38	pub(crate) extra: Vec<u8>,
39	pub(crate) row_count: u32,
40}
41
42#[instrument(
43	name = "wire::encode_frames",
44	level = "trace",
45	skip_all,
46	fields(
47		frame_count = frames.len(),
48		total_rows = frames.iter().map(|f| f.columns.first().map_or(0, |c| c.data.len())).sum::<usize>(),
49		bytes,
50	),
51)]
52pub fn encode_frames(frames: &[Frame], options: &EncodeOptions) -> Result<Vec<u8>, EncodeError> {
53	let mut buf = Vec::with_capacity(4096);
54	reserve_message_header(&mut buf);
55	for frame in frames {
56		encode_frame(frame, &mut buf, options)?;
57	}
58	write_message_header(&mut buf, frames.len() as u32);
59	Span::current().record("bytes", buf.len());
60	Ok(buf)
61}
62
63#[inline]
64fn reserve_message_header(buf: &mut Vec<u8>) {
65	buf.extend_from_slice(&[0u8; MESSAGE_HEADER_SIZE]);
66}
67
68#[inline]
69fn write_message_header(buf: &mut [u8], frame_count: u32) {
70	let total_size = buf.len() as u32;
71	buf[0..4].copy_from_slice(&RBCF_MAGIC.to_le_bytes());
72	buf[4..6].copy_from_slice(&RBCF_VERSION.to_le_bytes());
73	buf[6..8].copy_from_slice(&0u16.to_le_bytes());
74	buf[8..12].copy_from_slice(&frame_count.to_le_bytes());
75	buf[12..16].copy_from_slice(&total_size.to_le_bytes());
76}
77
78fn encode_frame(frame: &Frame, buf: &mut Vec<u8>, options: &EncodeOptions) -> Result<(), EncodeError> {
79	let frame_start = buf.len();
80	let row_count = frame.columns.first().map_or(0, |c| c.data.len()) as u32;
81	let column_count = frame.columns.len() as u16;
82	let meta_flags = compute_meta_flags(frame);
83
84	reserve_frame_header(buf);
85	write_frame_metadata(frame, meta_flags, buf);
86	encode_frame_columns(frame, buf, options)?;
87
88	let frame_size = (buf.len() - frame_start) as u32;
89	write_frame_header(buf, frame_start, row_count, column_count, meta_flags, frame_size);
90	Ok(())
91}
92
93#[inline]
94fn compute_meta_flags(frame: &Frame) -> u8 {
95	let mut flags = 0u8;
96	if !frame.row_numbers().is_empty() {
97		flags |= META_HAS_ROW_NUMBERS;
98	}
99	if !frame.created_at().is_empty() {
100		flags |= META_HAS_CREATED_AT;
101	}
102	if !frame.updated_at().is_empty() {
103		flags |= META_HAS_UPDATED_AT;
104	}
105	if !frame.time().is_empty() {
106		flags |= META_HAS_TIME;
107	}
108	flags
109}
110
111#[inline]
112fn reserve_frame_header(buf: &mut Vec<u8>) {
113	buf.extend_from_slice(&[0u8; FRAME_HEADER_SIZE]);
114}
115
116#[inline]
117fn write_frame_metadata(frame: &Frame, meta_flags: u8, buf: &mut Vec<u8>) {
118	if meta_flags & META_HAS_ROW_NUMBERS != 0 {
119		for rn in frame.row_numbers() {
120			buf.extend_from_slice(rn.to_le_bytes().as_ref());
121		}
122	}
123	if meta_flags & META_HAS_CREATED_AT != 0 {
124		for dt in frame.created_at() {
125			buf.extend_from_slice(dt.to_le_bytes().as_ref());
126		}
127	}
128	if meta_flags & META_HAS_UPDATED_AT != 0 {
129		for dt in frame.updated_at() {
130			buf.extend_from_slice(dt.to_le_bytes().as_ref());
131		}
132	}
133	if meta_flags & META_HAS_TIME != 0 {
134		for dt in frame.time() {
135			buf.extend_from_slice(dt.to_le_bytes().as_ref());
136		}
137	}
138}
139
140#[inline]
141fn encode_frame_columns(frame: &Frame, buf: &mut Vec<u8>, options: &EncodeOptions) -> Result<(), EncodeError> {
142	for col in &frame.columns {
143		encode_column(&col.name, &col.data, buf, options)?;
144	}
145	Ok(())
146}
147
148#[inline]
149fn write_frame_header(
150	buf: &mut [u8],
151	frame_start: usize,
152	row_count: u32,
153	column_count: u16,
154	meta_flags: u8,
155	frame_size: u32,
156) {
157	let h = frame_start;
158	buf[h..h + 4].copy_from_slice(&row_count.to_le_bytes());
159	buf[h + 4..h + 6].copy_from_slice(&column_count.to_le_bytes());
160	buf[h + 6] = meta_flags;
161	buf[h + 7] = 0;
162	buf[h + 8..h + 12].copy_from_slice(&frame_size.to_le_bytes());
163}
164
165fn encode_column(
166	name: &str,
167	col_data: &FrameColumnData,
168	buf: &mut Vec<u8>,
169	options: &EncodeOptions,
170) -> Result<(), EncodeError> {
171	let desired = options.force_encoding.unwrap_or_else(|| choose_encoding(col_data, options.compression));
172	let enc = try_encode_with(col_data, desired)?;
173	write_column(name, &enc, buf);
174	Ok(())
175}
176
177fn try_encode_with(col_data: &FrameColumnData, desired: Encoding) -> Result<EncodedColumn, EncodeError> {
178	let (inner, nones, has_nones) = match col_data {
179		FrameColumnData::Option {
180			inner,
181			bitvec,
182		} => {
183			let bitmap = encode_bitvec(bitvec);
184			(inner.as_ref(), bitmap, true)
185		}
186		other => (other, vec![], false),
187	};
188
189	let row_count = inner.len() as u32;
190
191	let result = match desired {
192		Encoding::Dict => varlen::try_dict_varlen(inner),
193		Encoding::Rle => fixed::try_rle_fixed(inner).or_else(|| varlen::try_rle_varlen(inner)),
194		Encoding::Delta => fixed::try_delta_fixed(inner),
195		Encoding::DeltaRle => fixed::try_delta_rle_fixed(inner),
196		_ => None,
197	};
198
199	let mut enc = match result {
200		Some(enc) => enc,
201		None => {
202			let plain = encode_plain(col_data)?;
203			EncodedColumn {
204				type_code: plain.type_code,
205				encoding: Encoding::Plain,
206				flags: 0,
207				nones: plain.nones,
208				data: plain.data,
209				offsets: plain.offsets,
210				extra: vec![],
211				row_count,
212			}
213		}
214	};
215
216	if has_nones {
217		enc.nones = nones;
218		enc.flags |= COL_FLAG_HAS_NONES;
219	}
220	enc.row_count = row_count;
221	Ok(enc)
222}
223
224fn write_column(name: &str, enc: &EncodedColumn, buf: &mut Vec<u8>) {
225	let name_bytes = name.as_bytes();
226	let name_len = name_bytes.len() as u16;
227	let name_pad = (4 - (name_bytes.len() % 4)) % 4;
228
229	buf.push(enc.type_code);
230	buf.push(enc.encoding as u8);
231	buf.push(enc.flags);
232	buf.push(0);
233	buf.extend_from_slice(&name_len.to_le_bytes());
234	buf.extend_from_slice(&0u16.to_le_bytes());
235	buf.extend_from_slice(&enc.row_count.to_le_bytes());
236	buf.extend_from_slice(&(enc.nones.len() as u32).to_le_bytes());
237	buf.extend_from_slice(&(enc.data.len() as u32).to_le_bytes());
238	buf.extend_from_slice(&(enc.offsets.len() as u32).to_le_bytes());
239	buf.extend_from_slice(&(enc.extra.len() as u32).to_le_bytes());
240
241	buf.extend_from_slice(name_bytes);
242	for _ in 0..name_pad {
243		buf.push(0);
244	}
245
246	buf.extend_from_slice(&enc.nones);
247
248	buf.extend_from_slice(&enc.data);
249
250	buf.extend_from_slice(&enc.offsets);
251
252	buf.extend_from_slice(&enc.extra);
253}