Skip to main content

micromegas_analytics/
metadata.rs

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};
27/// Type alias for shared, pre-serialized JSONB data.
28/// This represents JSONB properties that have been serialized once and can be reused.
29pub type SharedJsonbSerialized = Arc<Vec<u8>>;
30
31/// Analytics-optimized process metadata.
32///
33/// This struct is designed for analytics use cases where process properties need to be
34/// efficiently serialized to JSONB format multiple times. Unlike `ProcessInfo`, which
35/// uses `HashMap<String, String>` for properties, this struct stores pre-serialized
36/// JSONB data to eliminate redundant serialization overhead.
37#[derive(Debug, Clone)]
38pub struct ProcessMetadata {
39    // Core fields (same as ProcessInfo)
40    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/// Analytics-optimized stream metadata.
55///
56/// This struct is designed for analytics use cases where stream properties need to be
57/// efficiently serialized to JSONB format multiple times. Unlike `StreamInfo`, which
58/// uses `HashMap<String, String>` for properties, this struct stores pre-serialized
59/// JSONB data to eliminate redundant serialization overhead.
60#[derive(Debug, Clone)]
61pub struct StreamMetadata {
62    // Core fields (same as StreamInfo)
63    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    /// Creates StreamMetadata from StreamInfo by converting properties to JSONB format.
73    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
87/// Returns the thread name associated with the stream, if available.
88/// This function is only meaningful for streams associated with CPU threads.
89/// Returns format: "thread-name-thread-id" (e.g., "main-12345") or just "thread-id" if no name.
90pub 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    // thread_id falls back to stream_id
111    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    // Return "name-id" if name exists, otherwise just "id"
120    match thread_name {
121        Some(name) => Ok(format!("{}-{}", name, thread_id)),
122        None => Ok(thread_id),
123    }
124}
125
126/// Creates a `StreamMetadata` from a row of the `streams` view's (unprefixed) columns.
127#[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/// Finds a stream and its metadata using DataFusion (reads the `streams` view).
173#[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/// Creates a `ProcessMetadata` from a database row with pre-serialized JSONB properties.
217#[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/// Finds a process by its ID and returns it as ProcessMetadata with pre-serialized JSONB properties.
241#[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/// Finds a process and its latest timing information using DataFusion (optimized version).
273/// Returns (ProcessMetadata, last_block_end_ticks, last_block_end_time)
274#[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    // Extract all the required columns
320    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    // Handle properties column using PropertiesColumnAccessor
343    let properties_accessor = properties_column_by_name(batch, "properties")
344        .with_context(|| "accessing properties column")?;
345
346    // Get JSONB bytes directly from the properties column
347    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/// Creates a `BlockMetadata` from a recordbatch row.
375#[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}