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}