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