1use anyhow::{Context, Result};
2use chrono::{DateTime, Utc};
3use datafusion::arrow::array::{
4 BinaryArray, GenericListArray, Int32Array, Int64Array, RecordBatch, StringArray,
5 TimestampNanosecondArray,
6};
7use micromegas_telemetry::{
8 property::Property, stream_info::StreamInfo, types::block::BlockMetadata,
9};
10use micromegas_tracing::prelude::*;
11use micromegas_transit::{UserDefinedType, uuid_utils::parse_optional_uuid};
12use sqlx::Row;
13use std::sync::Arc;
14use uuid::Uuid;
15
16use crate::{
17 arrow_properties::serialize_properties_to_jsonb,
18 dfext::{string_column_accessor::string_column_by_name, typed_column::typed_column_by_name},
19 lakehouse::{
20 lakehouse_context::LakehouseContext, partition_cache::LivePartitionProvider,
21 query::make_session_context, session_configurator::NoOpSessionConfigurator,
22 view_factory::ViewFactory,
23 },
24 properties::properties_column_accessor::properties_column_by_name,
25 time::TimeRange,
26};
27pub type SharedJsonbSerialized = Arc<Vec<u8>>;
30
31#[derive(Debug, Clone)]
38pub struct ProcessMetadata {
39 pub process_id: uuid::Uuid,
41 pub exe: String,
42 pub username: String,
43 pub realname: String,
44 pub computer: String,
45 pub distro: String,
46 pub cpu_brand: String,
47 pub tsc_frequency: i64,
48 pub start_time: chrono::DateTime<chrono::Utc>,
49 pub start_ticks: i64,
50 pub parent_process_id: Option<uuid::Uuid>,
51 pub properties: SharedJsonbSerialized,
52}
53
54#[derive(Debug, Clone)]
61pub struct StreamMetadata {
62 pub process_id: Uuid,
64 pub stream_id: Uuid,
65 pub dependencies_metadata: Vec<UserDefinedType>,
66 pub objects_metadata: Vec<UserDefinedType>,
67 pub tags: Vec<String>,
68 pub properties: SharedJsonbSerialized,
69}
70
71impl StreamMetadata {
72 pub fn from_stream_info(stream_info: &StreamInfo) -> Result<Self> {
74 let properties = serialize_properties_to_jsonb(&stream_info.properties)
75 .with_context(|| "serializing stream properties to JSONB")?;
76 Ok(Self {
77 process_id: stream_info.process_id,
78 stream_id: stream_info.stream_id,
79 dependencies_metadata: stream_info.dependencies_metadata.clone(),
80 objects_metadata: stream_info.objects_metadata.clone(),
81 tags: stream_info.tags.clone(),
82 properties: Arc::new(properties),
83 })
84 }
85}
86
87pub fn get_thread_name_from_stream_metadata(stream: &StreamMetadata) -> Result<String> {
91 use jsonb::RawJsonb;
92
93 const THREAD_NAME_KEY: &str = "thread-name";
94 const THREAD_ID_KEY: &str = "thread-id";
95
96 if stream.properties.is_empty() {
97 return Ok(format!("{}", stream.stream_id));
98 }
99
100 let jsonb = RawJsonb::new(&stream.properties);
101
102 let thread_name = if let Ok(Some(thread_name_value)) = jsonb.get_by_name(THREAD_NAME_KEY, false)
103 && let Ok(Some(name)) = thread_name_value.as_raw().as_str()
104 {
105 Some(name.to_string())
106 } else {
107 None
108 };
109
110 let thread_id = if let Ok(Some(thread_id_value)) = jsonb.get_by_name(THREAD_ID_KEY, false)
112 && let Ok(Some(id)) = thread_id_value.as_raw().as_str()
113 {
114 id.to_string()
115 } else {
116 stream.stream_id.to_string()
117 };
118
119 match thread_name {
121 Some(name) => Ok(format!("{}-{}", name, thread_id)),
122 None => Ok(thread_id),
123 }
124}
125
126#[span_fn]
128pub fn stream_metadata_from_batch_row(batch: &RecordBatch, row: usize) -> Result<StreamMetadata> {
129 let stream_id_column = string_column_by_name(batch, "stream_id")?;
130 let process_id_column = string_column_by_name(batch, "process_id")?;
131 let dependencies_metadata_column: &BinaryArray =
132 typed_column_by_name(batch, "dependencies_metadata")?;
133 let objects_metadata_column: &BinaryArray = typed_column_by_name(batch, "objects_metadata")?;
134 let tags_column: &GenericListArray<i32> = typed_column_by_name(batch, "tags")?;
135 let properties_accessor = properties_column_by_name(batch, "properties")
136 .with_context(|| "accessing properties column")?;
137
138 let stream_id =
139 Uuid::parse_str(stream_id_column.value(row)?).with_context(|| "parsing stream_id")?;
140 let process_id =
141 Uuid::parse_str(process_id_column.value(row)?).with_context(|| "parsing process_id")?;
142
143 let dependencies_metadata: Vec<UserDefinedType> =
144 ciborium::from_reader(dependencies_metadata_column.value(row))
145 .with_context(|| "decoding dependencies_metadata")?;
146 let objects_metadata: Vec<UserDefinedType> =
147 ciborium::from_reader(objects_metadata_column.value(row))
148 .with_context(|| "decoding objects_metadata")?;
149 let tags: Vec<String> = tags_column
150 .value(row)
151 .as_any()
152 .downcast_ref::<StringArray>()
153 .with_context(|| "casting tags")?
154 .iter()
155 .map(|item| String::from(item.unwrap_or_default()))
156 .collect();
157
158 let properties = properties_accessor
159 .jsonb_value(row)
160 .with_context(|| "extracting JSONB from properties column")?;
161
162 Ok(StreamMetadata {
163 stream_id,
164 process_id,
165 dependencies_metadata,
166 objects_metadata,
167 tags,
168 properties: Arc::new(properties),
169 })
170}
171
172#[span_fn]
174pub async fn find_stream_from_view(
175 lakehouse: Arc<LakehouseContext>,
176 view_factory: Arc<ViewFactory>,
177 stream_id: &Uuid,
178 query_range: Option<TimeRange>,
179) -> Result<StreamMetadata> {
180 let partition_provider = Arc::new(LivePartitionProvider::new(lakehouse.lake().db_pool.clone()));
181
182 let ctx = make_session_context(
183 lakehouse,
184 partition_provider,
185 query_range,
186 view_factory,
187 Arc::new(NoOpSessionConfigurator),
188 false,
189 )
190 .await
191 .with_context(|| "creating DataFusion session context")?;
192
193 let sql = format!(
194 "SELECT stream_id, process_id, dependencies_metadata, objects_metadata, tags, properties
195 FROM streams
196 WHERE stream_id = '{stream_id}'"
197 );
198
199 let df = ctx
200 .sql(&sql)
201 .await
202 .with_context(|| "executing SQL query for stream")?;
203
204 let results = df
205 .collect()
206 .await
207 .with_context(|| "collecting results from DataFusion")?;
208
209 if results.is_empty() || results[0].num_rows() == 0 {
210 anyhow::bail!("Stream not found");
211 }
212
213 stream_metadata_from_batch_row(&results[0], 0)
214}
215
216#[span_fn]
218pub fn process_metadata_from_row(row: &sqlx::postgres::PgRow) -> Result<ProcessMetadata> {
219 let properties: Vec<Property> = row.try_get("process_properties")?;
220 let properties_map = micromegas_telemetry::property::into_hashmap(properties);
221 let serialized_properties = serialize_properties_to_jsonb(&properties_map)
222 .with_context(|| "serializing process properties to JSONB")?;
223
224 Ok(ProcessMetadata {
225 process_id: row.try_get("process_id")?,
226 exe: row.try_get("exe")?,
227 username: row.try_get("username")?,
228 realname: row.try_get("realname")?,
229 computer: row.try_get("computer")?,
230 distro: row.try_get("distro")?,
231 cpu_brand: row.try_get("cpu_brand")?,
232 tsc_frequency: row.try_get("tsc_frequency")?,
233 start_time: row.try_get("start_time")?,
234 start_ticks: row.try_get("start_ticks")?,
235 parent_process_id: row.try_get("parent_process_id")?,
236 properties: Arc::new(serialized_properties),
237 })
238}
239
240#[span_fn]
242pub async fn find_process(
243 pool: &sqlx::Pool<sqlx::Postgres>,
244 process_id: &sqlx::types::Uuid,
245) -> Result<ProcessMetadata> {
246 let row = instrument_named!(
247 sqlx::query(
248 "SELECT process_id,
249 exe,
250 username,
251 realname,
252 computer,
253 distro,
254 cpu_brand,
255 tsc_frequency,
256 start_time,
257 start_ticks,
258 parent_process_id,
259 properties as process_properties
260 FROM processes
261 WHERE process_id = $1;",
262 )
263 .bind(process_id)
264 .fetch_one(pool),
265 "sql_select_process"
266 )
267 .await
268 .with_context(|| "select from processes")?;
269 process_metadata_from_row(&row)
270}
271
272#[span_fn]
275pub async fn find_process_with_latest_timing(
276 lakehouse: Arc<LakehouseContext>,
277 view_factory: Arc<ViewFactory>,
278 process_id: &Uuid,
279 query_range: Option<TimeRange>,
280) -> Result<(ProcessMetadata, i64, DateTime<Utc>)> {
281 let partition_provider = Arc::new(LivePartitionProvider::new(lakehouse.lake().db_pool.clone()));
282
283 let ctx = make_session_context(
284 lakehouse,
285 partition_provider,
286 query_range,
287 view_factory,
288 Arc::new(NoOpSessionConfigurator),
289 false,
290 )
291 .await
292 .with_context(|| "creating DataFusion session context")?;
293
294 let sql = format!(
295 "SELECT process_id, exe, username, realname, computer, distro, cpu_brand,
296 tsc_frequency, start_time, start_ticks, parent_process_id, properties,
297 last_block_end_ticks, last_block_end_time
298 FROM processes
299 WHERE process_id = '{}'",
300 process_id
301 );
302
303 let df = ctx
304 .sql(&sql)
305 .await
306 .with_context(|| "executing SQL query for process with timing")?;
307
308 let results = df
309 .collect()
310 .await
311 .with_context(|| "collecting results from DataFusion")?;
312
313 if results.is_empty() || results[0].num_rows() == 0 {
314 anyhow::bail!("Process not found");
315 }
316
317 let batch = &results[0];
318
319 let process_id_column = string_column_by_name(batch, "process_id")?;
321 let exe_column = string_column_by_name(batch, "exe")?;
322 let username_column = string_column_by_name(batch, "username")?;
323 let realname_column = string_column_by_name(batch, "realname")?;
324 let computer_column = string_column_by_name(batch, "computer")?;
325 let distro_column = string_column_by_name(batch, "distro")?;
326 let cpu_brand_column = string_column_by_name(batch, "cpu_brand")?;
327 let tsc_frequency_column: &Int64Array = typed_column_by_name(batch, "tsc_frequency")?;
328 let start_time_column: &TimestampNanosecondArray = typed_column_by_name(batch, "start_time")?;
329 let start_ticks_column: &Int64Array = typed_column_by_name(batch, "start_ticks")?;
330 let last_block_end_ticks_column: &Int64Array =
331 typed_column_by_name(batch, "last_block_end_ticks")?;
332 let last_block_end_time_column: &TimestampNanosecondArray =
333 typed_column_by_name(batch, "last_block_end_time")?;
334 let parent_process_id_column = string_column_by_name(batch, "parent_process_id")?;
335
336 let parent_process_id = if parent_process_id_column.is_null(0) {
337 None
338 } else {
339 parse_optional_uuid(parent_process_id_column.value(0)?)?
340 };
341
342 let properties_accessor = properties_column_by_name(batch, "properties")
344 .with_context(|| "accessing properties column")?;
345
346 let properties_jsonb = Arc::new(
348 properties_accessor
349 .jsonb_value(0)
350 .with_context(|| "extracting JSONB from properties column")?,
351 );
352
353 let process_metadata = ProcessMetadata {
354 process_id: parse_optional_uuid(process_id_column.value(0)?)?
355 .ok_or_else(|| anyhow::anyhow!("process_id cannot be empty"))?,
356 exe: exe_column.value(0)?.to_string(),
357 username: username_column.value(0)?.to_string(),
358 realname: realname_column.value(0)?.to_string(),
359 computer: computer_column.value(0)?.to_string(),
360 distro: distro_column.value(0)?.to_string(),
361 cpu_brand: cpu_brand_column.value(0)?.to_string(),
362 tsc_frequency: tsc_frequency_column.value(0),
363 start_time: DateTime::from_timestamp_nanos(start_time_column.value(0)),
364 start_ticks: start_ticks_column.value(0),
365 parent_process_id,
366 properties: properties_jsonb,
367 };
368
369 let last_block_end_ticks = last_block_end_ticks_column.value(0);
370 let last_block_end_time = DateTime::from_timestamp_nanos(last_block_end_time_column.value(0));
371
372 Ok((process_metadata, last_block_end_ticks, last_block_end_time))
373}
374#[span_fn]
376pub fn block_from_batch_row(rb: &RecordBatch, row: usize) -> Result<BlockMetadata> {
377 let block_id_column = string_column_by_name(rb, "block_id")?;
378 let stream_id_column = string_column_by_name(rb, "stream_id")?;
379 let process_id_column = string_column_by_name(rb, "process_id")?;
380 let begin_time_column: &TimestampNanosecondArray = typed_column_by_name(rb, "begin_time")?;
381 let begin_ticks_column: &Int64Array = typed_column_by_name(rb, "begin_ticks")?;
382 let end_time_column: &TimestampNanosecondArray = typed_column_by_name(rb, "end_time")?;
383 let end_ticks_column: &Int64Array = typed_column_by_name(rb, "end_ticks")?;
384 let nb_objects_column: &Int32Array = typed_column_by_name(rb, "nb_objects")?;
385 let object_offset_column: &Int64Array = typed_column_by_name(rb, "object_offset")?;
386 let payload_size_column: &Int64Array = typed_column_by_name(rb, "payload_size")?;
387 let insert_time_column: &TimestampNanosecondArray = typed_column_by_name(rb, "insert_time")?;
388 Ok(BlockMetadata {
389 block_id: Uuid::parse_str(block_id_column.value(row)?)?,
390 stream_id: Uuid::parse_str(stream_id_column.value(row)?)?,
391 process_id: Uuid::parse_str(process_id_column.value(row)?)?,
392 begin_time: DateTime::from_timestamp_nanos(begin_time_column.value(row)),
393 end_time: DateTime::from_timestamp_nanos(end_time_column.value(row)),
394 begin_ticks: begin_ticks_column.value(row),
395 end_ticks: end_ticks_column.value(row),
396 nb_objects: nb_objects_column.value(row),
397 object_offset: object_offset_column.value(row),
398 payload_size: payload_size_column.value(row),
399 insert_time: DateTime::from_timestamp_nanos(insert_time_column.value(row)),
400 })
401}