1use core::time::Duration;
2use std::io;
3use std::sync::Arc;
4use std::time::SystemTime;
5
6use arrow::record_batch::RecordBatch;
7
8use arrow::array::{
9 ArrayRef, DictionaryArray, Int8Array, StringArray, TimestampMicrosecondArray, UInt8Array,
10 UInt32Array, UInt64Array,
11};
12use arrow::datatypes::{DataType, Field, Schema, TimeUnit};
13
14#[repr(u8)]
15#[derive(Clone, Copy)]
16pub enum Method {
17 Store = 0,
18 Deflate = 8,
19}
20
21pub struct ZipItemMeta {
22 pub name: String,
23 pub comment: String,
24 pub method: Method,
25 pub modified: SystemTime,
26 pub crc32: u32,
27 pub compressed_size: u64,
28 pub uncompressed_size: u64,
29}
30
31pub struct ZipFileMeta {
32 pub comment: String,
33 pub files: Vec<ZipItemMeta>,
34}
35
36pub struct Zip {
37 pub zip_id: String,
38 pub meta: ZipFileMeta,
39}
40
41pub fn stime2unixtime(s: SystemTime) -> Option<Duration> {
60 s.duration_since(SystemTime::UNIX_EPOCH).ok()
61}
62
63pub fn duration2us(d: Duration) -> Option<u64> {
64 let micros = d.as_micros();
65 micros.try_into().ok()
66}
67
68pub fn stime2unixtime_us(s: SystemTime) -> Option<u64> {
69 stime2unixtime(s).and_then(duration2us)
70}
71
72pub fn zip2record_batch(z: Zip) -> Result<RecordBatch, io::Error> {
73 let num_rows: usize = z.meta.files.len();
74
75 let single_zip_id = StringArray::from(vec![z.zip_id.clone()]);
76 let zip_id_indices = Int8Array::from(vec![0; num_rows]);
77 let zip_id_array: DictionaryArray<_> =
78 DictionaryArray::try_new(zip_id_indices, Arc::new(single_zip_id))
79 .map_err(io::Error::other)?;
80
81 let mut item_names: Vec<String> = Vec::with_capacity(num_rows);
82 let mut item_comments: Vec<String> = Vec::with_capacity(num_rows);
83 let mut item_modified: Vec<Option<i64>> = Vec::with_capacity(num_rows);
84 let mut crc32s: Vec<u32> = Vec::with_capacity(num_rows);
85 let mut compressed_sizes: Vec<u64> = Vec::with_capacity(num_rows);
86 let mut uncompressed_sizes: Vec<u64> = Vec::with_capacity(num_rows);
87
88 let mut methods: Vec<u8> = Vec::with_capacity(num_rows);
89
90 for f in &z.meta.files {
91 item_names.push(f.name.clone());
92 item_comments.push(f.comment.clone());
93
94 let us_opt: Option<i64> = stime2unixtime_us(f.modified).map(|u| u as i64);
95 item_modified.push(us_opt);
96
97 crc32s.push(f.crc32);
98 compressed_sizes.push(f.compressed_size);
99 uncompressed_sizes.push(f.uncompressed_size);
100
101 let m: Method = f.method;
102 let mu: u8 = m as u8;
103 methods.push(mu);
104 }
105
106 let schema = Schema::new(vec![
107 Field::new(
108 "zip_id",
109 DataType::Dictionary(Box::new(DataType::Int8), Box::new(DataType::Utf8)),
110 false,
111 ),
112 Field::new("item_name", DataType::Utf8, false),
113 Field::new("item_comment", DataType::Utf8, false),
114 Field::new("item_method", DataType::UInt8, false),
115 Field::new(
116 "item_modified",
117 DataType::Timestamp(TimeUnit::Microsecond, None),
118 true,
119 ),
120 Field::new("item_crc32", DataType::UInt32, false),
121 Field::new("item_compressed_size", DataType::UInt64, false),
122 Field::new("item_uncompressed_size", DataType::UInt64, false),
123 ]);
124
125 let arrays: Vec<ArrayRef> = vec![
126 Arc::new(zip_id_array),
127 Arc::new(StringArray::from(item_names)),
128 Arc::new(StringArray::from(item_comments)),
129 Arc::new(UInt8Array::from(methods)),
130 Arc::new(TimestampMicrosecondArray::from(item_modified)),
131 Arc::new(UInt32Array::from(crc32s)),
132 Arc::new(UInt64Array::from(compressed_sizes)),
133 Arc::new(UInt64Array::from(uncompressed_sizes)),
134 ];
135
136 RecordBatch::try_new(Arc::new(schema), arrays).map_err(io::Error::other)
137}