Skip to main content

reifydb_codec/frame/encode/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4pub(crate) mod any;
5mod fixed;
6mod varlen;
7
8use reifydb_value::{
9	encoding::LeBytes,
10	value::{
11		diff_type::DiffType,
12		frame::{data::FrameColumnData, frame::Frame},
13	},
14};
15use tracing::{Span, instrument};
16
17use crate::{
18	error::EncodeError,
19	frame::{
20		encoding::plain::{encode_bitvec, encode_plain},
21		format::{
22			COL_FLAG_HAS_NONES, Encoding, FRAME_HEADER_SIZE, MESSAGE_HEADER_SIZE, META_HAS_CREATED_AT,
23			META_HAS_ROW_NUMBERS, META_HAS_TIME, META_HAS_UPDATED_AT, RBCF_MAGIC, RBCF_VERSION,
24		},
25		heuristics::choose_encoding,
26		options::EncodeOptions,
27	},
28};
29
30pub(crate) struct EncodedColumn {
31	pub(crate) type_code: u8,
32	pub(crate) encoding: Encoding,
33	pub(crate) flags: u8,
34	pub(crate) nones: Vec<u8>,
35	pub(crate) data: Vec<u8>,
36	pub(crate) offsets: Vec<u8>,
37	pub(crate) extra: Vec<u8>,
38	pub(crate) row_count: u32,
39}
40
41#[instrument(
42	name = "wire::encode_frames",
43	level = "trace",
44	skip_all,
45	fields(
46		frame_count = frames.len(),
47		total_rows = frames.iter().map(|f| f.columns.first().map_or(0, |c| c.data.len())).sum::<usize>(),
48		bytes,
49	),
50)]
51pub fn encode_frames(frames: &[Frame], options: &EncodeOptions) -> Result<Vec<u8>, EncodeError> {
52	let mut buf = Vec::with_capacity(4096);
53	reserve_message_header(&mut buf);
54	for frame in frames {
55		encode_frame(frame, &mut buf, options)?;
56	}
57	write_message_header(&mut buf, frames.len() as u32);
58	Span::current().record("bytes", buf.len());
59	Ok(buf)
60}
61
62#[inline]
63fn reserve_message_header(buf: &mut Vec<u8>) {
64	buf.extend_from_slice(&[0u8; MESSAGE_HEADER_SIZE]);
65}
66
67#[inline]
68fn write_message_header(buf: &mut [u8], frame_count: u32) {
69	let total_size = buf.len() as u32;
70	buf[0..4].copy_from_slice(&RBCF_MAGIC.to_le_bytes());
71	buf[4..6].copy_from_slice(&RBCF_VERSION.to_le_bytes());
72	buf[6..8].copy_from_slice(&0u16.to_le_bytes());
73	buf[8..12].copy_from_slice(&frame_count.to_le_bytes());
74	buf[12..16].copy_from_slice(&total_size.to_le_bytes());
75}
76
77fn encode_frame(frame: &Frame, buf: &mut Vec<u8>, options: &EncodeOptions) -> Result<(), EncodeError> {
78	let frame_start = buf.len();
79	let row_count = frame.columns.first().map_or(0, |c| c.data.len()) as u32;
80	let column_count = frame.columns.len() as u16;
81	let meta_flags = compute_meta_flags(frame);
82
83	reserve_frame_header(buf);
84	write_frame_metadata(frame, meta_flags, buf);
85	encode_frame_columns(frame, buf, options)?;
86
87	let frame_size = (buf.len() - frame_start) as u32;
88	let op = frame.op.map_or(0, DiffType::as_u8);
89	write_frame_header(buf, frame_start, row_count, column_count, meta_flags, op, 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	op: u8,
156	frame_size: u32,
157) {
158	let h = frame_start;
159	buf[h..h + 4].copy_from_slice(&row_count.to_le_bytes());
160	buf[h + 4..h + 6].copy_from_slice(&column_count.to_le_bytes());
161	buf[h + 6] = meta_flags;
162	buf[h + 7] = op;
163	buf[h + 8..h + 12].copy_from_slice(&frame_size.to_le_bytes());
164}
165
166fn encode_column(
167	name: &str,
168	col_data: &FrameColumnData,
169	buf: &mut Vec<u8>,
170	options: &EncodeOptions,
171) -> Result<(), EncodeError> {
172	let desired = options.force_encoding.unwrap_or_else(|| choose_encoding(col_data, options.compression));
173	let enc = try_encode_with(col_data, desired)?;
174	write_column(name, &enc, buf);
175	Ok(())
176}
177
178fn try_encode_with(col_data: &FrameColumnData, desired: Encoding) -> Result<EncodedColumn, EncodeError> {
179	let (inner, nones, has_nones) = match col_data {
180		FrameColumnData::Option {
181			inner,
182			bitvec,
183		} => {
184			let bitmap = encode_bitvec(bitvec);
185			(inner.as_ref(), bitmap, true)
186		}
187		other => (other, vec![], false),
188	};
189
190	let row_count = inner.len() as u32;
191
192	let result = match desired {
193		Encoding::Dict => varlen::try_dict_varlen(inner),
194		Encoding::Rle => fixed::try_rle_fixed(inner).or_else(|| varlen::try_rle_varlen(inner)),
195		Encoding::Delta => fixed::try_delta_fixed(inner),
196		Encoding::DeltaRle => fixed::try_delta_rle_fixed(inner),
197		_ => None,
198	};
199
200	let mut enc = match result {
201		Some(enc) => enc,
202		None => {
203			let plain = encode_plain(col_data)?;
204			EncodedColumn {
205				type_code: plain.type_code,
206				encoding: Encoding::Plain,
207				flags: 0,
208				nones: plain.nones,
209				data: plain.data,
210				offsets: plain.offsets,
211				extra: vec![],
212				row_count,
213			}
214		}
215	};
216
217	if has_nones {
218		enc.nones = nones;
219		enc.flags |= COL_FLAG_HAS_NONES;
220	}
221	enc.row_count = row_count;
222	Ok(enc)
223}
224
225fn write_column(name: &str, enc: &EncodedColumn, buf: &mut Vec<u8>) {
226	let name_bytes = name.as_bytes();
227	let name_len = name_bytes.len() as u16;
228	let name_pad = (4 - (name_bytes.len() % 4)) % 4;
229
230	buf.push(enc.type_code);
231	buf.push(enc.encoding as u8);
232	buf.push(enc.flags);
233	buf.push(0);
234	buf.extend_from_slice(&name_len.to_le_bytes());
235	buf.extend_from_slice(&0u16.to_le_bytes());
236	buf.extend_from_slice(&enc.row_count.to_le_bytes());
237	buf.extend_from_slice(&(enc.nones.len() as u32).to_le_bytes());
238	buf.extend_from_slice(&(enc.data.len() as u32).to_le_bytes());
239	buf.extend_from_slice(&(enc.offsets.len() as u32).to_le_bytes());
240	buf.extend_from_slice(&(enc.extra.len() as u32).to_le_bytes());
241
242	buf.extend_from_slice(name_bytes);
243	for _ in 0..name_pad {
244		buf.push(0);
245	}
246
247	buf.extend_from_slice(&enc.nones);
248
249	buf.extend_from_slice(&enc.data);
250
251	buf.extend_from_slice(&enc.offsets);
252
253	buf.extend_from_slice(&enc.extra);
254}