Skip to main content

micromegas_analytics/
log_entries_table.rs

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
23/// Returns the schema for the log entries table.
24pub 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
86/// A builder for creating a `RecordBatch` of log entries.
87pub 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        // assuming that the events are in order
127        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    /// Append only per-entry variable data (optimized for batch processing)
161    pub fn append_entry_only(&mut self, row: &LogEntry) -> Result<()> {
162        // Only append fields that truly vary per log entry
163        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    /// Batch fill all constant columns for all entries in block
172    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        // For PrimitiveBuilder (insert_times): use append_slice for better performance
183        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}