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