1use std::sync::Arc;
2
3use anyhow::{Context, Result};
4use chrono::DateTime;
5use datafusion::arrow::array::ArrayBuilder;
6use datafusion::arrow::array::BinaryDictionaryBuilder;
7use datafusion::arrow::array::PrimitiveBuilder;
8use datafusion::arrow::array::StringBuilder;
9use datafusion::arrow::array::StringDictionaryBuilder;
10use datafusion::arrow::datatypes::DataType;
11use datafusion::arrow::datatypes::Field;
12use datafusion::arrow::datatypes::Int32Type;
13use datafusion::arrow::datatypes::Schema;
14use datafusion::arrow::datatypes::TimeUnit;
15use datafusion::arrow::datatypes::TimestampNanosecondType;
16use datafusion::arrow::record_batch::RecordBatch;
17
18use crate::log_entry::LogEntry;
19use crate::metadata::ProcessMetadata;
20use crate::properties::property_set_jsonb_dictionary_builder::PropertySetJsonbDictionaryBuilder;
21use crate::time::TimeRange;
22
23pub fn log_table_schema() -> Schema {
25 Schema::new(vec![
26 Field::new(
27 "process_id",
28 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
29 false,
30 ),
31 Field::new(
32 "stream_id",
33 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
34 false,
35 ),
36 Field::new(
37 "block_id",
38 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
39 false,
40 ),
41 Field::new(
42 "insert_time",
43 DataType::Timestamp(TimeUnit::Nanosecond, Some("+00:00".into())),
44 false,
45 ),
46 Field::new(
47 "exe",
48 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
49 false,
50 ),
51 Field::new(
52 "username",
53 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
54 false,
55 ),
56 Field::new(
57 "computer",
58 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
59 false,
60 ),
61 Field::new(
62 "time",
63 DataType::Timestamp(TimeUnit::Nanosecond, Some("+00:00".into())),
64 false,
65 ),
66 Field::new(
67 "target",
68 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8)),
69 false,
70 ),
71 Field::new("level", DataType::Int32, false),
72 Field::new("msg", DataType::Utf8, false),
73 Field::new(
74 "properties",
75 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Binary)),
76 false,
77 ),
78 Field::new(
79 "process_properties",
80 DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Binary)),
81 false,
82 ),
83 ])
84}
85
86pub struct LogEntriesRecordBuilder {
88 process_ids: StringDictionaryBuilder<Int32Type>,
89 stream_ids: StringDictionaryBuilder<Int32Type>,
90 block_ids: StringDictionaryBuilder<Int32Type>,
91 insert_times: PrimitiveBuilder<TimestampNanosecondType>,
92 exes: StringDictionaryBuilder<Int32Type>,
93 usernames: StringDictionaryBuilder<Int32Type>,
94 computers: StringDictionaryBuilder<Int32Type>,
95 times: PrimitiveBuilder<TimestampNanosecondType>,
96 targets: StringDictionaryBuilder<Int32Type>,
97 levels: PrimitiveBuilder<Int32Type>,
98 msgs: StringBuilder,
99 properties: PropertySetJsonbDictionaryBuilder,
100 process_properties: BinaryDictionaryBuilder<Int32Type>,
101}
102
103impl LogEntriesRecordBuilder {
104 pub fn with_capacity(capacity: usize) -> Self {
105 Self {
106 process_ids: StringDictionaryBuilder::new(),
107 stream_ids: StringDictionaryBuilder::new(),
108 block_ids: StringDictionaryBuilder::new(),
109 insert_times: PrimitiveBuilder::with_capacity(capacity),
110 exes: StringDictionaryBuilder::new(),
111 usernames: StringDictionaryBuilder::new(),
112 computers: StringDictionaryBuilder::new(),
113 times: PrimitiveBuilder::with_capacity(capacity),
114 targets: StringDictionaryBuilder::new(),
115 levels: PrimitiveBuilder::with_capacity(capacity),
116 msgs: StringBuilder::new(),
117 properties: PropertySetJsonbDictionaryBuilder::new(capacity),
118 process_properties: BinaryDictionaryBuilder::new(),
119 }
120 }
121
122 pub fn get_time_range(&self) -> Option<TimeRange> {
123 if self.is_empty() {
124 return None;
125 }
126 let slice = self.times.values_slice();
128 Some(TimeRange::new(
129 DateTime::from_timestamp_nanos(slice[0]),
130 DateTime::from_timestamp_nanos(slice[slice.len() - 1]),
131 ))
132 }
133
134 pub fn len(&self) -> i64 {
135 self.times.len() as i64
136 }
137
138 pub fn is_empty(&self) -> bool {
139 self.times.len() == 0
140 }
141
142 pub fn append(&mut self, row: &LogEntry) -> Result<()> {
143 self.process_ids
144 .append(format!("{}", row.process.process_id))?;
145 self.stream_ids.append(&*row.stream_id)?;
146 self.block_ids.append(&*row.block_id)?;
147 self.insert_times.append_value(row.insert_time);
148 self.exes.append(&row.process.exe)?;
149 self.usernames.append(&row.process.username)?;
150 self.computers.append(&row.process.computer)?;
151 self.times.append_value(row.time);
152 self.targets.append(row.target)?;
153 self.levels.append_value(row.level);
154 self.msgs.append_value(row.msg);
155 self.properties.append_property_set(&row.properties)?;
156 self.process_properties.append(&*row.process.properties)?;
157 Ok(())
158 }
159
160 pub fn append_entry_only(&mut self, row: &LogEntry) -> Result<()> {
162 self.times.append_value(row.time);
164 self.targets.append(row.target)?;
165 self.levels.append_value(row.level);
166 self.msgs.append_value(row.msg);
167 self.properties.append_property_set(&row.properties)?;
168 Ok(())
169 }
170
171 pub fn fill_constant_columns(
173 &mut self,
174 process: &ProcessMetadata,
175 stream_id: &str,
176 block_id: &str,
177 insert_time: i64,
178 entry_count: usize,
179 ) -> Result<()> {
180 let process_id_str = format!("{}", process.process_id);
181
182 let insert_times_slice = vec![insert_time; entry_count];
184 self.insert_times.append_slice(&insert_times_slice);
185
186 self.process_ids.append_n(&process_id_str, entry_count)?;
187 self.stream_ids.append_n(stream_id, entry_count)?;
188 self.block_ids.append_n(block_id, entry_count)?;
189 self.exes.append_n(&process.exe, entry_count)?;
190 self.usernames.append_n(&process.username, entry_count)?;
191 self.computers.append_n(&process.computer, entry_count)?;
192 self.process_properties
193 .append_n(&**process.properties, entry_count)?;
194
195 Ok(())
196 }
197
198 pub fn finish(mut self) -> Result<RecordBatch> {
199 RecordBatch::try_new(
200 Arc::new(log_table_schema()),
201 vec![
202 Arc::new(self.process_ids.finish()),
203 Arc::new(self.stream_ids.finish()),
204 Arc::new(self.block_ids.finish()),
205 Arc::new(self.insert_times.finish().with_timezone_utc()),
206 Arc::new(self.exes.finish()),
207 Arc::new(self.usernames.finish()),
208 Arc::new(self.computers.finish()),
209 Arc::new(self.times.finish().with_timezone_utc()),
210 Arc::new(self.targets.finish()),
211 Arc::new(self.levels.finish()),
212 Arc::new(self.msgs.finish()),
213 Arc::new(self.properties.finish()?),
214 Arc::new(self.process_properties.finish()),
215 ],
216 )
217 .with_context(|| "building record batch")
218 }
219}