Skip to main content

rs_zip2meta2rbat/
core.rs

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
41/*
42pub struct FlatZipMetaInfo {
43    pub zip_id: String,
44    pub item_name: String,
45    pub item_comment: String,
46    pub item_modified: SystemTime,
47    pub crc32: u32,
48    pub compressed_size: u64,
49    pub uncompressed_size: u64,
50}
51
52impl FlatZipMetaInfo {
53    pub fn item_modified_unixtime_us(&self) -> Option<u64> {
54        stime2unixtime_us(self.item_modified)
55    }
56}
57*/
58
59pub 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}