Skip to main content

micromegas_telemetry_sink/
stream_block.rs

1use anyhow::Result;
2use micromegas_telemetry::{block_wire_format, compression::compress, wire_format::encode_cbor};
3use micromegas_tracing::{
4    event::{EventBlock, ExtractDeps, TracingBlock},
5    images::ImageBlock,
6    logs::LogBlock,
7    metrics::MetricsBlock,
8    prelude::*,
9    spans::ThreadBlock,
10};
11use micromegas_transit::HeterogeneousQueue;
12
13pub trait StreamBlock: Send + Sync {
14    /// Encodes the stream block into a binary format.
15    ///
16    /// This function serializes the block data, compresses it, and then encodes it
17    /// into the wire format for transmission.
18    ///
19    /// # Arguments
20    ///
21    /// * `process_info` - Information about the current process, used for time calibration.
22    fn encode_bin(&self, process_info: &ProcessInfo) -> Result<Vec<u8>>;
23}
24
25fn encode_block<Q>(block: &EventBlock<Q>, process_info: &ProcessInfo) -> Result<Vec<u8>>
26where
27    Q: HeterogeneousQueue + ExtractDeps,
28    <Q as ExtractDeps>::DepsQueue: HeterogeneousQueue,
29{
30    let block_id = uuid::Uuid::new_v4();
31    trace!("encoding block_id={block_id}");
32    let end = block.end.as_ref().unwrap();
33
34    let payload = block_wire_format::BlockPayload {
35        dependencies: compress(block.events.extract().as_bytes())?,
36        objects: compress(block.events.as_bytes())?,
37    };
38
39    let block = block_wire_format::Block {
40        block_id,
41        stream_id: block.stream_id,
42        process_id: block.process_id,
43        begin_time: block
44            .begin
45            .time
46            .to_rfc3339_opts(chrono::SecondsFormat::Nanos, false),
47        begin_ticks: block.begin.ticks - process_info.start_ticks,
48        end_time: end
49            .time
50            .to_rfc3339_opts(chrono::SecondsFormat::Nanos, false),
51        end_ticks: end.ticks - process_info.start_ticks,
52        payload,
53        nb_objects: block.nb_objects() as i32,
54        object_offset: block.object_offset() as i64,
55    };
56    encode_cbor(&block)
57}
58
59impl StreamBlock for ImageBlock {
60    fn encode_bin(&self, process_info: &ProcessInfo) -> Result<Vec<u8>> {
61        encode_block(self, process_info)
62    }
63}
64
65impl StreamBlock for LogBlock {
66    fn encode_bin(&self, process_info: &ProcessInfo) -> Result<Vec<u8>> {
67        encode_block(self, process_info)
68    }
69}
70
71impl StreamBlock for MetricsBlock {
72    fn encode_bin(&self, process_info: &ProcessInfo) -> Result<Vec<u8>> {
73        encode_block(self, process_info)
74    }
75}
76
77impl StreamBlock for ThreadBlock {
78    fn encode_bin(&self, process_info: &ProcessInfo) -> Result<Vec<u8>> {
79        encode_block(self, process_info)
80    }
81}