Skip to main content

micromegas_analytics/
images_table.rs

1use anyhow::{Context, Result};
2use chrono::DateTime;
3use datafusion::arrow::{
4    array::{
5        ArrayBuilder, BinaryBuilder, PrimitiveBuilder, StringBuilder, StringDictionaryBuilder,
6    },
7    datatypes::{DataType, Field, Int32Type, Int64Type, Schema, TimeUnit, TimestampNanosecondType},
8    record_batch::RecordBatch,
9};
10use std::sync::Arc;
11
12use crate::{metadata::ProcessMetadata, time::TimeRange};
13
14pub fn images_table_schema() -> Schema {
15    Schema::new(vec![
16        Field::new(
17            "process_id",
18            DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
19            false,
20        ),
21        Field::new(
22            "stream_id",
23            DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
24            false,
25        ),
26        Field::new(
27            "block_id",
28            DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
29            false,
30        ),
31        Field::new(
32            "insert_time",
33            DataType::Timestamp(TimeUnit::Nanosecond, Some("+00:00".into())),
34            false,
35        ),
36        Field::new(
37            "exe",
38            DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
39            false,
40        ),
41        Field::new(
42            "username",
43            DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
44            false,
45        ),
46        Field::new(
47            "computer",
48            DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
49            false,
50        ),
51        Field::new(
52            "time",
53            DataType::Timestamp(TimeUnit::Nanosecond, Some("+00:00".into())),
54            false,
55        ),
56        Field::new("name", DataType::Utf8, false),
57        Field::new(
58            "format",
59            DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
60            false,
61        ),
62        Field::new("payload_size", DataType::Int64, false),
63        Field::new("data", DataType::Binary, false),
64    ])
65}
66
67pub struct ImagesRecordBuilder {
68    process_ids: StringDictionaryBuilder<Int32Type>,
69    stream_ids: StringDictionaryBuilder<Int32Type>,
70    block_ids: StringDictionaryBuilder<Int32Type>,
71    insert_times: PrimitiveBuilder<TimestampNanosecondType>,
72    exes: StringDictionaryBuilder<Int32Type>,
73    usernames: StringDictionaryBuilder<Int32Type>,
74    computers: StringDictionaryBuilder<Int32Type>,
75    times: PrimitiveBuilder<TimestampNanosecondType>,
76    names: StringBuilder,
77    formats: StringDictionaryBuilder<Int32Type>,
78    payload_sizes: PrimitiveBuilder<Int64Type>,
79    data: BinaryBuilder,
80    min_time: Option<i64>,
81    max_time: Option<i64>,
82}
83
84impl Default for ImagesRecordBuilder {
85    fn default() -> Self {
86        Self::new()
87    }
88}
89
90impl ImagesRecordBuilder {
91    pub fn new() -> Self {
92        Self {
93            process_ids: StringDictionaryBuilder::new(),
94            stream_ids: StringDictionaryBuilder::new(),
95            block_ids: StringDictionaryBuilder::new(),
96            insert_times: PrimitiveBuilder::new(),
97            exes: StringDictionaryBuilder::new(),
98            usernames: StringDictionaryBuilder::new(),
99            computers: StringDictionaryBuilder::new(),
100            times: PrimitiveBuilder::new(),
101            names: StringBuilder::new(),
102            formats: StringDictionaryBuilder::new(),
103            payload_sizes: PrimitiveBuilder::new(),
104            data: BinaryBuilder::new(),
105            min_time: None,
106            max_time: None,
107        }
108    }
109
110    pub fn is_empty(&self) -> bool {
111        self.times.len() == 0
112    }
113
114    pub fn get_time_range(&self) -> Option<TimeRange> {
115        match (self.min_time, self.max_time) {
116            (Some(min_ns), Some(max_ns)) => Some(TimeRange::new(
117                DateTime::from_timestamp_nanos(min_ns),
118                DateTime::from_timestamp_nanos(max_ns),
119            )),
120            _ => None,
121        }
122    }
123
124    #[allow(clippy::too_many_arguments)]
125    pub fn append(
126        &mut self,
127        process: &ProcessMetadata,
128        process_id_str: &str,
129        stream_id_str: &str,
130        block_id_str: &str,
131        insert_time_nanos: i64,
132        time_ns: i64,
133        name: &str,
134        format: &str,
135        payload_size: i64,
136        image_data: &[u8],
137    ) -> Result<()> {
138        self.process_ids.append(process_id_str)?;
139        self.stream_ids.append(stream_id_str)?;
140        self.block_ids.append(block_id_str)?;
141        self.insert_times.append_value(insert_time_nanos);
142        self.exes.append(&process.exe)?;
143        self.usernames.append(&process.username)?;
144        self.computers.append(&process.computer)?;
145        self.times.append_value(time_ns);
146        self.names.append_value(name);
147        self.formats.append(format)?;
148        self.payload_sizes.append_value(payload_size);
149        self.data.append_value(image_data);
150        self.min_time = Some(self.min_time.map(|m| m.min(time_ns)).unwrap_or(time_ns));
151        self.max_time = Some(self.max_time.map(|m| m.max(time_ns)).unwrap_or(time_ns));
152        Ok(())
153    }
154
155    pub fn finish(mut self) -> Result<RecordBatch> {
156        RecordBatch::try_new(
157            Arc::new(images_table_schema()),
158            vec![
159                Arc::new(self.process_ids.finish()),
160                Arc::new(self.stream_ids.finish()),
161                Arc::new(self.block_ids.finish()),
162                Arc::new(self.insert_times.finish().with_timezone_utc()),
163                Arc::new(self.exes.finish()),
164                Arc::new(self.usernames.finish()),
165                Arc::new(self.computers.finish()),
166                Arc::new(self.times.finish().with_timezone_utc()),
167                Arc::new(self.names.finish()),
168                Arc::new(self.formats.finish()),
169                Arc::new(self.payload_sizes.finish()),
170                Arc::new(self.data.finish()),
171            ],
172        )
173        .with_context(|| "building images record batch")
174    }
175}