1use crate::error::SimError;
22use crate::recorder::Recorder;
23
24pub const PARQUET_PAGE_ROWS: usize = 1024;
27
28const MAGIC: &[u8; 4] = b"PAR1";
30
31const TYPE_DOUBLE: i32 = 5;
33const REPETITION_REQUIRED: i32 = 0;
34const ENCODING_PLAIN: i32 = 0;
35const ENCODING_RLE: i32 = 3;
36const CODEC_UNCOMPRESSED: i32 = 0;
37const PAGE_DATA: i32 = 0;
38
39const T_I16: u8 = 4;
41const T_I32: u8 = 5;
42const T_I64: u8 = 6;
43const T_BINARY: u8 = 8;
44const T_LIST: u8 = 9;
45const T_STRUCT: u8 = 12;
46
47struct Compact {
49 out: Vec<u8>,
50 last_field: Vec<i16>,
52}
53
54impl Compact {
55 fn new() -> Self {
56 Self {
57 out: Vec::new(),
58 last_field: vec![0],
59 }
60 }
61
62 fn varint(&mut self, mut value: u64) {
64 while value >= 0x80 {
65 self.out.push((value as u8 & 0x7f) | 0x80);
67 value >>= 7;
68 }
69 self.out.push(value as u8);
70 }
71
72 fn int(&mut self, value: i64) {
74 self.varint(((value << 1) ^ (value >> 63)) as u64);
75 }
76
77 fn field(&mut self, id: i16, kind: u8) {
80 let last = self.last_field.last().copied().unwrap_or(0);
81 let delta = id - last;
82 if (1..=15).contains(&delta) {
83 self.out.push(((delta as u8) << 4) | kind);
84 } else {
85 self.out.push(kind);
86 self.int(i64::from(id));
87 }
88 if let Some(slot) = self.last_field.last_mut() {
89 *slot = id;
90 }
91 }
92
93 fn i16(&mut self, id: i16, value: i16) {
94 self.field(id, T_I16);
95 self.int(i64::from(value));
96 }
97
98 fn i32(&mut self, id: i16, value: i32) {
99 self.field(id, T_I32);
100 self.int(i64::from(value));
101 }
102
103 fn i64(&mut self, id: i16, value: i64) {
104 self.field(id, T_I64);
105 self.int(value);
106 }
107
108 fn bytes(&mut self, value: &[u8]) {
109 self.varint(value.len() as u64);
110 self.out.extend_from_slice(value);
111 }
112
113 fn string(&mut self, id: i16, value: &str) {
114 self.field(id, T_BINARY);
115 self.bytes(value.as_bytes());
116 }
117
118 fn list(&mut self, id: i16, element: u8, len: usize) {
121 self.field(id, T_LIST);
122 if len < 15 {
123 self.out.push(((len as u8) << 4) | element);
124 } else {
125 self.out.push(0xf0 | element);
126 self.varint(len as u64);
127 }
128 }
129
130 fn open(&mut self) {
132 self.last_field.push(0);
133 }
134
135 fn open_field(&mut self, id: i16) {
137 self.field(id, T_STRUCT);
138 self.open();
139 }
140
141 fn close(&mut self) {
143 self.out.push(0);
144 self.last_field.pop();
145 }
146}
147
148fn count<T: TryFrom<usize>>(value: usize) -> Result<T, SimError> {
150 T::try_from(value).map_err(|_| SimError::Unsupported {
151 what: "a recording too large for a Parquet file's counts",
152 })
153}
154
155const KEY_VALUE: [(&str, &str); 3] = [
158 ("tool", hpr_core::tool::NAME),
159 ("tool_version", hpr_core::tool::VERSION),
160 ("designation", hpr_core::tool::DESIGNATION),
161];
162
163struct Chunk {
165 offset: usize,
166 size: usize,
167}
168
169pub fn parquet(recorder: &Recorder) -> Result<Vec<u8>, SimError> {
183 let columns = recorder.columns();
184 if columns.is_empty() {
185 return Err(SimError::Unsupported {
186 what: "a Parquet file with no columns",
187 });
188 }
189 let rows = recorder.rows();
190 for row in rows {
191 for &value in row {
192 super::finite("recorded value", value)?;
193 }
194 }
195
196 let mut out = MAGIC.to_vec();
197 let mut chunks = Vec::with_capacity(columns.len());
198 if !rows.is_empty() {
199 for column in 0..columns.len() {
200 let offset = out.len();
201 for page in rows.chunks(PARQUET_PAGE_ROWS) {
202 let size: i32 = count(page.len() * 8)?;
203 let mut header = Compact::new();
204 header.i32(1, PAGE_DATA);
205 header.i32(2, size);
206 header.i32(3, size);
207 header.open_field(5);
208 header.i32(1, count(page.len())?);
209 header.i32(2, ENCODING_PLAIN);
210 header.i32(3, ENCODING_RLE);
211 header.i32(4, ENCODING_RLE);
212 header.close();
213 header.close();
214 out.extend_from_slice(&header.out);
215 for row in page {
216 out.extend_from_slice(&row[column].to_le_bytes());
217 }
218 }
219 chunks.push(Chunk {
220 offset,
221 size: out.len() - offset,
222 });
223 }
224 }
225
226 let mut footer = Compact::new();
227 footer.i32(1, 1);
228 footer.list(2, T_STRUCT, columns.len() + 1);
229 footer.open();
230 footer.string(4, "schema");
231 footer.i32(5, count(columns.len())?);
232 footer.close();
233 for name in &columns {
234 footer.open();
235 footer.i32(1, TYPE_DOUBLE);
236 footer.i32(3, REPETITION_REQUIRED);
237 footer.string(4, name);
238 footer.close();
239 }
240 footer.i64(3, count(rows.len())?);
241 footer.list(4, T_STRUCT, usize::from(!chunks.is_empty()));
242 if !chunks.is_empty() {
243 let total: i64 = count(chunks.iter().map(|c| c.size).sum())?;
244 footer.open();
245 footer.list(1, T_STRUCT, chunks.len());
246 for (name, chunk) in columns.iter().zip(&chunks) {
247 let size: i64 = count(chunk.size)?;
248 footer.open();
249 footer.i64(2, 0);
251 footer.open_field(3);
252 footer.i32(1, TYPE_DOUBLE);
253 footer.list(2, T_I32, 1);
254 footer.int(i64::from(ENCODING_PLAIN));
255 footer.list(3, T_BINARY, 1);
256 footer.bytes(name.as_bytes());
257 footer.i32(4, CODEC_UNCOMPRESSED);
258 footer.i64(5, count(rows.len())?);
259 footer.i64(6, size);
260 footer.i64(7, size);
261 footer.i64(9, count(chunk.offset)?);
262 footer.close();
263 footer.close();
264 }
265 footer.i64(2, total);
266 footer.i64(3, count(rows.len())?);
267 footer.i64(5, count(MAGIC.len())?);
268 footer.i64(6, total);
269 footer.i16(7, 0);
270 footer.close();
271 }
272 footer.list(5, T_STRUCT, KEY_VALUE.len());
275 for (key, value) in KEY_VALUE {
276 footer.open();
277 footer.string(1, key);
278 footer.string(2, value);
279 footer.close();
280 }
281 footer.string(6, concat!("hpr-sim version ", env!("CARGO_PKG_VERSION")));
282 footer.close();
283
284 let length: u32 = count(footer.out.len())?;
285 out.extend_from_slice(&footer.out);
286 out.extend_from_slice(&length.to_le_bytes());
287 out.extend_from_slice(MAGIC);
288 Ok(out)
289}
290
291#[cfg(test)]
292mod tests {
293 use bytes::Bytes;
294 use parquet_reader::basic::{Encoding, PageType, Repetition, Type};
295 use parquet_reader::file::reader::{FileReader, SerializedFileReader};
296 use parquet_reader::record::RowAccessor;
297
298 use super::*;
299 use crate::flight::{FlightSettings, Simulation};
300 use crate::rail::Rail;
301 use crate::recorder::Channel;
302 use hpr_atmos::ConstantWind;
303
304 use crate::testing::{design, windy_environment};
305
306 fn read(file: Vec<u8>) -> (SerializedFileReader<Bytes>, Vec<String>, Vec<Vec<f64>>) {
308 let reader = SerializedFileReader::new(Bytes::from(file)).unwrap();
309 let schema = reader.metadata().file_metadata().schema_descr_ptr();
310 let names = schema
311 .columns()
312 .iter()
313 .map(|c| c.name().to_owned())
314 .collect();
315 for column in schema.columns() {
316 assert_eq!(column.physical_type(), Type::DOUBLE);
317 assert_eq!(
318 column.self_type().get_basic_info().repetition(),
319 Repetition::REQUIRED
320 );
321 }
322 let rows = reader
323 .get_row_iter(None)
324 .unwrap()
325 .map(|row| {
326 let row = row.unwrap();
327 (0..row.len()).map(|i| row.get_double(i).unwrap()).collect()
328 })
329 .collect();
330 (reader, names, rows)
331 }
332
333 #[test]
334 fn apaches_reader_reads_back_every_recorded_value_exactly() {
335 let (recorder, ..) = super::super::tests::flown();
336 assert!(recorder.rows().len() > 20);
337 let (reader, names, rows) = read(parquet(&recorder).unwrap());
338 assert_eq!(names, recorder.columns());
339 assert_eq!(rows, recorder.rows());
340 let metadata = reader.metadata();
341 assert_eq!(metadata.num_row_groups(), 1);
342 assert_eq!(
343 metadata.file_metadata().num_rows(),
344 recorder.rows().len() as i64
345 );
346 assert_eq!(
347 metadata.file_metadata().created_by(),
348 Some(concat!("hpr-sim version ", env!("CARGO_PKG_VERSION")))
349 );
350 }
351
352 #[test]
355 fn the_key_value_metadata_names_the_program_version_and_designation() {
356 let (recorder, ..) = super::super::tests::flown();
357 let empty = Recorder::with_rows(vec![Channel::Time], Vec::new());
358 for recorder in [recorder, empty] {
359 let (reader, names, rows) = read(parquet(&recorder).unwrap());
360 assert_eq!(names, recorder.columns());
361 assert_eq!(rows, recorder.rows());
362 let pairs: Vec<(String, Option<String>)> = reader
363 .metadata()
364 .file_metadata()
365 .key_value_metadata()
366 .unwrap()
367 .iter()
368 .map(|kv| (kv.key.clone(), kv.value.clone()))
369 .collect();
370 let pair = |key: &str, value: &str| (key.to_owned(), Some(value.to_owned()));
371 assert_eq!(
372 pairs,
373 [
374 pair("tool", "hpr-sim"),
375 pair("tool_version", env!("CARGO_PKG_VERSION")),
376 pair("designation", "FS · SW · TOOL 005"),
377 ]
378 );
379 }
380 }
381
382 #[test]
385 fn a_long_recording_of_every_channel_splits_into_pages() {
386 let sim = Simulation::new(
387 &design("rocketpy-valetudo"),
388 "example",
389 windy_environment(ConstantWind::new(4.0, 0.0).unwrap()),
390 Rail::vertical(3.0),
391 FlightSettings::default(),
392 )
393 .unwrap();
394 let mut recorder = Recorder::new(Channel::ALL.to_vec(), Some(0.01)).unwrap();
395 sim.run(&mut recorder).unwrap();
396 let n = recorder.rows().len();
397 assert!(n > 2 * PARQUET_PAGE_ROWS, "{n} rows");
398 assert!(recorder.columns().len() >= 15);
399
400 let file = parquet(&recorder).unwrap();
401 let (reader, names, rows) = read(file.clone());
402 assert_eq!(names, recorder.columns());
403 assert_eq!(rows, recorder.rows());
404 let metadata = reader.metadata().row_group(0);
407 let mut end = MAGIC.len() as i64;
408 for chunk in metadata.columns() {
409 assert_eq!(chunk.data_page_offset(), end);
410 assert_eq!(chunk.compressed_size(), chunk.uncompressed_size());
411 assert_eq!(chunk.num_values(), n as i64);
412 end += chunk.compressed_size();
413 }
414 let length = u32::from_le_bytes(file[file.len() - 8..file.len() - 4].try_into().unwrap());
415 assert_eq!(end, (file.len() - 8 - length as usize) as i64);
416 assert_eq!(metadata.total_byte_size(), end - MAGIC.len() as i64);
417 assert_eq!(metadata.compressed_size(), metadata.total_byte_size());
418 assert_eq!(metadata.file_offset(), Some(MAGIC.len() as i64));
419 assert_eq!(metadata.ordinal(), Some(0));
420
421 let group = reader.get_row_group(0).unwrap();
422 for column in 0..names.len() {
423 let pages: Vec<_> = group
424 .get_column_page_reader(column)
425 .unwrap()
426 .map(Result::unwrap)
427 .collect();
428 assert_eq!(pages.len(), n.div_ceil(PARQUET_PAGE_ROWS));
429 let mut values = 0;
430 for page in &pages {
431 assert_eq!(page.page_type(), PageType::DATA_PAGE);
432 assert_eq!(page.encoding(), Encoding::PLAIN);
433 assert!(page.num_values() as usize <= PARQUET_PAGE_ROWS);
434 values += page.num_values() as usize;
435 }
436 assert_eq!(values, n);
437 }
438 }
439
440 #[test]
441 fn an_empty_recording_is_a_file_with_no_row_groups() {
442 let recorder = Recorder::with_rows(vec![Channel::Time, Channel::Mach], Vec::new());
443 let (reader, names, rows) = read(parquet(&recorder).unwrap());
444 assert_eq!(names, ["time_s", "mach"]);
445 assert!(rows.is_empty());
446 assert_eq!(reader.metadata().num_row_groups(), 0);
447 assert_eq!(reader.metadata().file_metadata().num_rows(), 0);
448 }
449
450 #[test]
451 fn extreme_values_keep_their_bits() {
452 let values = [
453 -0.0,
454 f64::MIN_POSITIVE,
455 5e-324,
456 f64::MAX,
457 f64::MIN,
458 0.1 + 0.2,
459 ];
460 let recorder = Recorder::with_rows(
461 vec![Channel::Time, Channel::Mach],
462 values.iter().map(|&v| vec![v, -v]).collect(),
463 );
464 let (_, _, rows) = read(parquet(&recorder).unwrap());
465 assert_eq!(rows.len(), values.len());
466 for (row, &v) in rows.iter().zip(&values) {
467 assert_eq!(row[0].to_bits(), v.to_bits());
468 assert_eq!(row[1].to_bits(), (-v).to_bits());
469 }
470 }
471
472 #[test]
473 fn no_columns_or_a_value_that_is_not_finite_is_refused() {
474 let empty = Recorder::with_rows(Vec::new(), vec![Vec::new()]);
475 assert!(matches!(
476 parquet(&empty),
477 Err(SimError::Unsupported { what }) if what.contains("no columns")
478 ));
479 for bad in [f64::NAN, f64::INFINITY, f64::NEG_INFINITY] {
480 let recorder = Recorder::with_rows(vec![Channel::Time], vec![vec![0.0], vec![bad]]);
481 assert!(matches!(
482 parquet(&recorder),
483 Err(SimError::Domain { what: "recorded value", value })
484 if value.to_bits() == bad.to_bits()
485 ));
486 }
487 }
488
489 #[test]
493 fn a_small_file_is_the_bytes_the_specification_gives() {
494 let recorder = Recorder::with_rows(vec![Channel::Time], vec![vec![1.0], vec![2.0]]);
495 let mut expected: Vec<u8> = b"PAR1".to_vec();
496 expected.extend([
501 0x15, 0x00, 0x15, 0x20, 0x15, 0x20, 0x2c, 0x15, 0x04, 0x15, 0x00, 0x15, 0x06, 0x15,
502 0x06, 0x00, 0x00,
503 ]);
504 expected.extend(1.0_f64.to_le_bytes());
506 expected.extend(2.0_f64.to_le_bytes());
507 let footer_start = expected.len();
508 expected.extend([0x15, 0x02, 0x19, 0x2c, 0x48, 0x06]);
511 expected.extend(b"schema");
512 expected.extend([0x15, 0x02, 0x00, 0x15, 0x0a, 0x25, 0x00, 0x18, 0x06]);
513 expected.extend(b"time_s");
514 expected.extend([0x00, 0x16, 0x04]);
515 expected.extend([
519 0x19, 0x1c, 0x19, 0x1c, 0x26, 0x00, 0x1c, 0x15, 0x0a, 0x19, 0x15,
520 ]);
521 expected.extend([0x00, 0x19, 0x18, 0x06]);
522 expected.extend(b"time_s");
523 expected.extend([
524 0x15, 0x00, 0x16, 0x04, 0x16, 0x42, 0x16, 0x42, 0x26, 0x08, 0x00, 0x00,
525 ]);
526 expected.extend([
529 0x16, 0x42, 0x16, 0x04, 0x26, 0x08, 0x16, 0x42, 0x14, 0x00, 0x00,
530 ]);
531 let version = env!("CARGO_PKG_VERSION");
538 expected.extend([0x19, 0x3c]);
539 expected.extend([0x18, 0x04]);
540 expected.extend(b"tool");
541 expected.extend([0x18, 0x07]);
542 expected.extend(b"hpr-sim");
543 expected.push(0x00);
544 expected.extend([0x18, 0x0c]);
545 expected.extend(b"tool_version");
546 expected.extend([0x18, version.len() as u8]);
547 expected.extend(version.as_bytes());
548 expected.push(0x00);
549 expected.extend([0x18, 0x0b]);
550 expected.extend(b"designation");
551 expected.extend([0x18, 0x14]);
552 expected.extend(b"FS \xc2\xb7 SW \xc2\xb7 TOOL 005");
553 expected.push(0x00);
554 let created_by = concat!("hpr-sim version ", env!("CARGO_PKG_VERSION"));
556 expected.extend([0x18, created_by.len() as u8]);
557 expected.extend(created_by.as_bytes());
558 expected.push(0x00);
559 let footer_length = (expected.len() - footer_start) as u32;
560 expected.extend(footer_length.to_le_bytes());
561 expected.extend(b"PAR1");
562
563 assert_eq!(parquet(&recorder).unwrap(), expected);
564 }
565
566 #[test]
568 fn compact_protocol_matches_its_specification() {
569 let mut c = Compact::new();
570 c.varint(50399);
572 assert_eq!(c.out, [0xdf, 0x89, 0x03]);
573 for (value, expected) in [
575 (0, 0u64),
576 (-1, 1),
577 (1, 2),
578 (-2, 3),
579 (2_147_483_647, 4_294_967_294),
580 (-2_147_483_648, 4_294_967_295),
581 ] {
582 let mut c = Compact::new();
583 c.int(value);
584 let mut v = Compact::new();
585 v.varint(expected);
586 assert_eq!(c.out, v.out, "{value}");
587 }
588 let mut c = Compact::new();
591 c.i32(1, 7);
592 c.i32(17, 7);
593 assert_eq!(c.out, [0x15, 0x0e, 0x05, 0x22, 0x0e]);
594 let mut c = Compact::new();
596 c.list(1, T_I32, 14);
597 c.list(2, T_STRUCT, 15);
598 assert_eq!(c.out, [0x19, 0xe5, 0x19, 0xfc, 0x0f]);
599 let mut c = Compact::new();
601 c.open_field(3);
602 c.i16(1, -1);
603 c.close();
604 c.i64(4, 1);
605 assert_eq!(c.out, [0x3c, 0x14, 0x01, 0x00, 0x16, 0x02]);
606 }
607}