reifydb_codec/frame/encode/
mod.rs1pub(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}