micromegas_telemetry_sink/
stream_block.rs1use 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 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}