1use byteorder::{LittleEndian, ReadBytesExt, WriteBytesExt};
7use std::collections::HashMap;
8use std::io::{Read, Write};
9
10use crate::error::{IoError, Result};
11
12use super::types::{ColumnData, EncodingType};
13
14pub fn write_plain_f64<W: Write>(writer: &mut W, data: &[f64]) -> Result<()> {
20 for &val in data {
21 writer
22 .write_f64::<LittleEndian>(val)
23 .map_err(|e| IoError::FileError(format!("Failed to write f64: {}", e)))?;
24 }
25 Ok(())
26}
27
28pub fn read_plain_f64<R: Read>(reader: &mut R, count: usize) -> Result<Vec<f64>> {
30 let mut data = Vec::with_capacity(count);
31 for _ in 0..count {
32 let val = reader
33 .read_f64::<LittleEndian>()
34 .map_err(|e| IoError::FormatError(format!("Failed to read f64: {}", e)))?;
35 data.push(val);
36 }
37 Ok(data)
38}
39
40pub fn write_plain_i64<W: Write>(writer: &mut W, data: &[i64]) -> Result<()> {
42 for &val in data {
43 writer
44 .write_i64::<LittleEndian>(val)
45 .map_err(|e| IoError::FileError(format!("Failed to write i64: {}", e)))?;
46 }
47 Ok(())
48}
49
50pub fn read_plain_i64<R: Read>(reader: &mut R, count: usize) -> Result<Vec<i64>> {
52 let mut data = Vec::with_capacity(count);
53 for _ in 0..count {
54 let val = reader
55 .read_i64::<LittleEndian>()
56 .map_err(|e| IoError::FormatError(format!("Failed to read i64: {}", e)))?;
57 data.push(val);
58 }
59 Ok(data)
60}
61
62pub fn write_plain_str<W: Write>(writer: &mut W, data: &[String]) -> Result<()> {
64 for s in data {
65 let bytes = s.as_bytes();
66 writer
67 .write_u32::<LittleEndian>(bytes.len() as u32)
68 .map_err(|e| IoError::FileError(format!("Failed to write string length: {}", e)))?;
69 writer
70 .write_all(bytes)
71 .map_err(|e| IoError::FileError(format!("Failed to write string data: {}", e)))?;
72 }
73 Ok(())
74}
75
76pub fn read_plain_str<R: Read>(reader: &mut R, count: usize) -> Result<Vec<String>> {
78 let mut data = Vec::with_capacity(count);
79 for _ in 0..count {
80 let len = reader
81 .read_u32::<LittleEndian>()
82 .map_err(|e| IoError::FormatError(format!("Failed to read string length: {}", e)))?
83 as usize;
84 let mut buf = vec![0u8; len];
85 reader
86 .read_exact(&mut buf)
87 .map_err(|e| IoError::FormatError(format!("Failed to read string data: {}", e)))?;
88 let s = String::from_utf8(buf)
89 .map_err(|e| IoError::FormatError(format!("Invalid UTF-8 string: {}", e)))?;
90 data.push(s);
91 }
92 Ok(data)
93}
94
95pub fn write_plain_bool<W: Write>(writer: &mut W, data: &[bool]) -> Result<()> {
97 let num_bytes = (data.len() + 7) / 8;
99 let mut packed = vec![0u8; num_bytes];
100 for (i, &val) in data.iter().enumerate() {
101 if val {
102 packed[i / 8] |= 1 << (i % 8);
103 }
104 }
105 writer
106 .write_all(&packed)
107 .map_err(|e| IoError::FileError(format!("Failed to write bool data: {}", e)))?;
108 Ok(())
109}
110
111pub fn read_plain_bool<R: Read>(reader: &mut R, count: usize) -> Result<Vec<bool>> {
113 let num_bytes = (count + 7) / 8;
114 let mut packed = vec![0u8; num_bytes];
115 reader
116 .read_exact(&mut packed)
117 .map_err(|e| IoError::FormatError(format!("Failed to read bool data: {}", e)))?;
118 let mut data = Vec::with_capacity(count);
119 for i in 0..count {
120 data.push((packed[i / 8] >> (i % 8)) & 1 == 1);
121 }
122 Ok(data)
123}
124
125pub fn write_rle_f64<W: Write>(writer: &mut W, data: &[f64]) -> Result<()> {
131 if data.is_empty() {
132 return Ok(());
133 }
134 let mut i = 0;
135 while i < data.len() {
136 let val = data[i];
137 let mut run_len: u32 = 1;
138 while i + (run_len as usize) < data.len() && data[i + (run_len as usize)] == val {
139 run_len += 1;
140 }
141 writer
142 .write_u32::<LittleEndian>(run_len)
143 .map_err(|e| IoError::FileError(format!("Failed to write RLE run length: {}", e)))?;
144 writer
145 .write_f64::<LittleEndian>(val)
146 .map_err(|e| IoError::FileError(format!("Failed to write RLE value: {}", e)))?;
147 i += run_len as usize;
148 }
149 Ok(())
150}
151
152pub fn read_rle_f64<R: Read>(reader: &mut R, total_count: usize) -> Result<Vec<f64>> {
154 let mut data = Vec::with_capacity(total_count);
155 while data.len() < total_count {
156 let run_len = reader
157 .read_u32::<LittleEndian>()
158 .map_err(|e| IoError::FormatError(format!("Failed to read RLE run length: {}", e)))?
159 as usize;
160 let val = reader
161 .read_f64::<LittleEndian>()
162 .map_err(|e| IoError::FormatError(format!("Failed to read RLE value: {}", e)))?;
163 for _ in 0..run_len {
164 data.push(val);
165 }
166 }
167 if data.len() != total_count {
168 return Err(IoError::FormatError(format!(
169 "RLE decoded {} values, expected {}",
170 data.len(),
171 total_count
172 )));
173 }
174 Ok(data)
175}
176
177pub fn write_rle_i64<W: Write>(writer: &mut W, data: &[i64]) -> Result<()> {
179 if data.is_empty() {
180 return Ok(());
181 }
182 let mut i = 0;
183 while i < data.len() {
184 let val = data[i];
185 let mut run_len: u32 = 1;
186 while i + (run_len as usize) < data.len() && data[i + (run_len as usize)] == val {
187 run_len += 1;
188 }
189 writer
190 .write_u32::<LittleEndian>(run_len)
191 .map_err(|e| IoError::FileError(format!("Failed to write RLE run length: {}", e)))?;
192 writer
193 .write_i64::<LittleEndian>(val)
194 .map_err(|e| IoError::FileError(format!("Failed to write RLE value: {}", e)))?;
195 i += run_len as usize;
196 }
197 Ok(())
198}
199
200pub fn read_rle_i64<R: Read>(reader: &mut R, total_count: usize) -> Result<Vec<i64>> {
202 let mut data = Vec::with_capacity(total_count);
203 while data.len() < total_count {
204 let run_len = reader
205 .read_u32::<LittleEndian>()
206 .map_err(|e| IoError::FormatError(format!("Failed to read RLE run length: {}", e)))?
207 as usize;
208 let val = reader
209 .read_i64::<LittleEndian>()
210 .map_err(|e| IoError::FormatError(format!("Failed to read RLE value: {}", e)))?;
211 for _ in 0..run_len {
212 data.push(val);
213 }
214 }
215 if data.len() != total_count {
216 return Err(IoError::FormatError(format!(
217 "RLE decoded {} values, expected {}",
218 data.len(),
219 total_count
220 )));
221 }
222 Ok(data)
223}
224
225pub fn write_rle_str<W: Write>(writer: &mut W, data: &[String]) -> Result<()> {
227 if data.is_empty() {
228 return Ok(());
229 }
230 let mut i = 0;
231 while i < data.len() {
232 let val = &data[i];
233 let mut run_len: u32 = 1;
234 while i + (run_len as usize) < data.len() && &data[i + (run_len as usize)] == val {
235 run_len += 1;
236 }
237 writer
238 .write_u32::<LittleEndian>(run_len)
239 .map_err(|e| IoError::FileError(format!("Failed to write RLE run length: {}", e)))?;
240 let bytes = val.as_bytes();
241 writer
242 .write_u32::<LittleEndian>(bytes.len() as u32)
243 .map_err(|e| IoError::FileError(format!("Failed to write RLE string length: {}", e)))?;
244 writer
245 .write_all(bytes)
246 .map_err(|e| IoError::FileError(format!("Failed to write RLE string data: {}", e)))?;
247 i += run_len as usize;
248 }
249 Ok(())
250}
251
252pub fn read_rle_str<R: Read>(reader: &mut R, total_count: usize) -> Result<Vec<String>> {
254 let mut data = Vec::with_capacity(total_count);
255 while data.len() < total_count {
256 let run_len = reader
257 .read_u32::<LittleEndian>()
258 .map_err(|e| IoError::FormatError(format!("Failed to read RLE run length: {}", e)))?
259 as usize;
260 let str_len = reader
261 .read_u32::<LittleEndian>()
262 .map_err(|e| IoError::FormatError(format!("Failed to read RLE string length: {}", e)))?
263 as usize;
264 let mut buf = vec![0u8; str_len];
265 reader
266 .read_exact(&mut buf)
267 .map_err(|e| IoError::FormatError(format!("Failed to read RLE string data: {}", e)))?;
268 let s = String::from_utf8(buf)
269 .map_err(|e| IoError::FormatError(format!("Invalid UTF-8 in RLE string: {}", e)))?;
270 for _ in 0..run_len {
271 data.push(s.clone());
272 }
273 }
274 if data.len() != total_count {
275 return Err(IoError::FormatError(format!(
276 "RLE decoded {} values, expected {}",
277 data.len(),
278 total_count
279 )));
280 }
281 Ok(data)
282}
283
284pub fn write_rle_bool<W: Write>(writer: &mut W, data: &[bool]) -> Result<()> {
286 if data.is_empty() {
287 return Ok(());
288 }
289 let mut i = 0;
290 while i < data.len() {
291 let val = data[i];
292 let mut run_len: u32 = 1;
293 while i + (run_len as usize) < data.len() && data[i + (run_len as usize)] == val {
294 run_len += 1;
295 }
296 writer
297 .write_u32::<LittleEndian>(run_len)
298 .map_err(|e| IoError::FileError(format!("Failed to write RLE run length: {}", e)))?;
299 writer
300 .write_u8(if val { 1 } else { 0 })
301 .map_err(|e| IoError::FileError(format!("Failed to write RLE bool value: {}", e)))?;
302 i += run_len as usize;
303 }
304 Ok(())
305}
306
307pub fn read_rle_bool<R: Read>(reader: &mut R, total_count: usize) -> Result<Vec<bool>> {
309 let mut data = Vec::with_capacity(total_count);
310 while data.len() < total_count {
311 let run_len = reader
312 .read_u32::<LittleEndian>()
313 .map_err(|e| IoError::FormatError(format!("Failed to read RLE run length: {}", e)))?
314 as usize;
315 let val = reader
316 .read_u8()
317 .map_err(|e| IoError::FormatError(format!("Failed to read RLE bool value: {}", e)))?
318 != 0;
319 for _ in 0..run_len {
320 data.push(val);
321 }
322 }
323 if data.len() != total_count {
324 return Err(IoError::FormatError(format!(
325 "RLE decoded {} values, expected {}",
326 data.len(),
327 total_count
328 )));
329 }
330 Ok(data)
331}
332
333pub fn write_dict_str<W: Write>(writer: &mut W, data: &[String]) -> Result<()> {
339 let mut dictionary: Vec<String> = Vec::new();
341 let mut dict_map: HashMap<String, u32> = HashMap::new();
342 let mut indices: Vec<u32> = Vec::with_capacity(data.len());
343
344 for s in data {
345 let idx = if let Some(&existing) = dict_map.get(s) {
346 existing
347 } else {
348 let new_idx = dictionary.len() as u32;
349 dict_map.insert(s.clone(), new_idx);
350 dictionary.push(s.clone());
351 new_idx
352 };
353 indices.push(idx);
354 }
355
356 writer
358 .write_u32::<LittleEndian>(dictionary.len() as u32)
359 .map_err(|e| IoError::FileError(format!("Failed to write dictionary size: {}", e)))?;
360
361 for entry in &dictionary {
363 let bytes = entry.as_bytes();
364 writer
365 .write_u32::<LittleEndian>(bytes.len() as u32)
366 .map_err(|e| IoError::FileError(format!("Failed to write dict entry length: {}", e)))?;
367 writer
368 .write_all(bytes)
369 .map_err(|e| IoError::FileError(format!("Failed to write dict entry: {}", e)))?;
370 }
371
372 for &idx in &indices {
374 writer
375 .write_u32::<LittleEndian>(idx)
376 .map_err(|e| IoError::FileError(format!("Failed to write dict index: {}", e)))?;
377 }
378
379 Ok(())
380}
381
382pub fn read_dict_str<R: Read>(reader: &mut R, count: usize) -> Result<Vec<String>> {
384 let dict_size = reader
386 .read_u32::<LittleEndian>()
387 .map_err(|e| IoError::FormatError(format!("Failed to read dictionary size: {}", e)))?
388 as usize;
389
390 let mut dictionary = Vec::with_capacity(dict_size);
391 for _ in 0..dict_size {
392 let len = reader
393 .read_u32::<LittleEndian>()
394 .map_err(|e| IoError::FormatError(format!("Failed to read dict entry length: {}", e)))?
395 as usize;
396 let mut buf = vec![0u8; len];
397 reader
398 .read_exact(&mut buf)
399 .map_err(|e| IoError::FormatError(format!("Failed to read dict entry: {}", e)))?;
400 let s = String::from_utf8(buf)
401 .map_err(|e| IoError::FormatError(format!("Invalid UTF-8 in dict entry: {}", e)))?;
402 dictionary.push(s);
403 }
404
405 let mut data = Vec::with_capacity(count);
407 for _ in 0..count {
408 let idx = reader
409 .read_u32::<LittleEndian>()
410 .map_err(|e| IoError::FormatError(format!("Failed to read dict index: {}", e)))?
411 as usize;
412 if idx >= dictionary.len() {
413 return Err(IoError::FormatError(format!(
414 "Dictionary index {} out of range (dict size {})",
415 idx,
416 dictionary.len()
417 )));
418 }
419 data.push(dictionary[idx].clone());
420 }
421
422 Ok(data)
423}
424
425pub fn write_delta_f64<W: Write>(writer: &mut W, data: &[f64]) -> Result<()> {
431 if data.is_empty() {
432 return Ok(());
433 }
434 writer
436 .write_f64::<LittleEndian>(data[0])
437 .map_err(|e| IoError::FileError(format!("Failed to write delta base: {}", e)))?;
438 for i in 1..data.len() {
440 let delta = data[i] - data[i - 1];
441 writer
442 .write_f64::<LittleEndian>(delta)
443 .map_err(|e| IoError::FileError(format!("Failed to write delta: {}", e)))?;
444 }
445 Ok(())
446}
447
448pub fn read_delta_f64<R: Read>(reader: &mut R, count: usize) -> Result<Vec<f64>> {
450 if count == 0 {
451 return Ok(Vec::new());
452 }
453 let mut data = Vec::with_capacity(count);
454 let base = reader
455 .read_f64::<LittleEndian>()
456 .map_err(|e| IoError::FormatError(format!("Failed to read delta base: {}", e)))?;
457 data.push(base);
458 for _ in 1..count {
459 let delta = reader
460 .read_f64::<LittleEndian>()
461 .map_err(|e| IoError::FormatError(format!("Failed to read delta: {}", e)))?;
462 let prev = data[data.len() - 1];
463 data.push(prev + delta);
464 }
465 Ok(data)
466}
467
468pub fn write_delta_i64<W: Write>(writer: &mut W, data: &[i64]) -> Result<()> {
470 if data.is_empty() {
471 return Ok(());
472 }
473 writer
474 .write_i64::<LittleEndian>(data[0])
475 .map_err(|e| IoError::FileError(format!("Failed to write delta base: {}", e)))?;
476 for i in 1..data.len() {
477 let delta = data[i].wrapping_sub(data[i - 1]);
478 writer
479 .write_i64::<LittleEndian>(delta)
480 .map_err(|e| IoError::FileError(format!("Failed to write delta: {}", e)))?;
481 }
482 Ok(())
483}
484
485pub fn read_delta_i64<R: Read>(reader: &mut R, count: usize) -> Result<Vec<i64>> {
487 if count == 0 {
488 return Ok(Vec::new());
489 }
490 let mut data = Vec::with_capacity(count);
491 let base = reader
492 .read_i64::<LittleEndian>()
493 .map_err(|e| IoError::FormatError(format!("Failed to read delta base: {}", e)))?;
494 data.push(base);
495 for _ in 1..count {
496 let delta = reader
497 .read_i64::<LittleEndian>()
498 .map_err(|e| IoError::FormatError(format!("Failed to read delta: {}", e)))?;
499 let prev = data[data.len() - 1];
500 data.push(prev.wrapping_add(delta));
501 }
502 Ok(data)
503}
504
505pub fn encode_column<W: Write>(
511 writer: &mut W,
512 data: &ColumnData,
513 encoding: EncodingType,
514) -> Result<()> {
515 match (data, encoding) {
516 (ColumnData::Float64(v), EncodingType::Plain) => write_plain_f64(writer, v),
518 (ColumnData::Int64(v), EncodingType::Plain) => write_plain_i64(writer, v),
519 (ColumnData::Str(v), EncodingType::Plain) => write_plain_str(writer, v),
520 (ColumnData::Bool(v), EncodingType::Plain) => write_plain_bool(writer, v),
521 (ColumnData::Float64(v), EncodingType::Rle) => write_rle_f64(writer, v),
523 (ColumnData::Int64(v), EncodingType::Rle) => write_rle_i64(writer, v),
524 (ColumnData::Str(v), EncodingType::Rle) => write_rle_str(writer, v),
525 (ColumnData::Bool(v), EncodingType::Rle) => write_rle_bool(writer, v),
526 (ColumnData::Str(v), EncodingType::Dictionary) => write_dict_str(writer, v),
528 (ColumnData::Float64(v), EncodingType::Delta) => write_delta_f64(writer, v),
530 (ColumnData::Int64(v), EncodingType::Delta) => write_delta_i64(writer, v),
531 (data, enc) => Err(IoError::FormatError(format!(
533 "Encoding {:?} not supported for column type {:?}",
534 enc,
535 data.type_tag()
536 ))),
537 }
538}
539
540pub fn decode_column<R: Read>(
542 reader: &mut R,
543 type_tag: super::types::ColumnTypeTag,
544 encoding: EncodingType,
545 count: usize,
546) -> Result<ColumnData> {
547 use super::types::ColumnTypeTag;
548
549 match (type_tag, encoding) {
550 (ColumnTypeTag::Float64, EncodingType::Plain) => {
552 Ok(ColumnData::Float64(read_plain_f64(reader, count)?))
553 }
554 (ColumnTypeTag::Int64, EncodingType::Plain) => {
555 Ok(ColumnData::Int64(read_plain_i64(reader, count)?))
556 }
557 (ColumnTypeTag::Str, EncodingType::Plain) => {
558 Ok(ColumnData::Str(read_plain_str(reader, count)?))
559 }
560 (ColumnTypeTag::Bool, EncodingType::Plain) => {
561 Ok(ColumnData::Bool(read_plain_bool(reader, count)?))
562 }
563 (ColumnTypeTag::Float64, EncodingType::Rle) => {
565 Ok(ColumnData::Float64(read_rle_f64(reader, count)?))
566 }
567 (ColumnTypeTag::Int64, EncodingType::Rle) => {
568 Ok(ColumnData::Int64(read_rle_i64(reader, count)?))
569 }
570 (ColumnTypeTag::Str, EncodingType::Rle) => {
571 Ok(ColumnData::Str(read_rle_str(reader, count)?))
572 }
573 (ColumnTypeTag::Bool, EncodingType::Rle) => {
574 Ok(ColumnData::Bool(read_rle_bool(reader, count)?))
575 }
576 (ColumnTypeTag::Str, EncodingType::Dictionary) => {
578 Ok(ColumnData::Str(read_dict_str(reader, count)?))
579 }
580 (ColumnTypeTag::Float64, EncodingType::Delta) => {
582 Ok(ColumnData::Float64(read_delta_f64(reader, count)?))
583 }
584 (ColumnTypeTag::Int64, EncodingType::Delta) => {
585 Ok(ColumnData::Int64(read_delta_i64(reader, count)?))
586 }
587 (tt, enc) => Err(IoError::FormatError(format!(
589 "Encoding {:?} not supported for type {:?}",
590 enc, tt
591 ))),
592 }
593}