Skip to main content

micromegas_analytics/
response_writer.rs

1use anyhow::{Context, Result};
2use async_trait::async_trait;
3use bytes::Bytes;
4use micromegas_telemetry::wire_format::encode_cbor;
5use micromegas_tracing::prelude::*;
6use tokio::sync::mpsc::Sender;
7
8/// A trait for writing log entries.
9#[async_trait]
10pub trait Logger: Send + Sync {
11    async fn write_log_entry(&self, msg: String) -> Result<()>;
12}
13
14/// A writer for sending responses to a client.
15pub struct ResponseWriter {
16    sender: Option<Sender<Bytes>>,
17}
18
19impl ResponseWriter {
20    pub fn new(sender: Option<Sender<Bytes>>) -> Self {
21        Self { sender }
22    }
23    pub async fn write_string(&self, value: String) -> Result<()> {
24        info!("{value}");
25        let buffer = encode_cbor(&value)?;
26        if let Some(sender) = &self.sender {
27            sender
28                .send(buffer.into())
29                .await
30                .with_context(|| "writing response")?;
31        }
32        Ok(())
33    }
34
35    pub fn is_closed(&self) -> bool {
36        if let Some(sender) = &self.sender {
37            sender.is_closed()
38        } else {
39            false
40        }
41    }
42}
43
44#[async_trait]
45impl Logger for ResponseWriter {
46    async fn write_log_entry(&self, msg: String) -> Result<()> {
47        self.write_string(msg).await
48    }
49}
50
51/// A sender for sending log entries to a channel.
52///
53/// The channel item is `Result<(time, msg), String>` so a paired consumer (`AsyncLogStream`) can
54/// also carry a query-level failure through the same channel; `write_log_entry` always sends
55/// `Ok(...)`, so this is behavior-preserving for every existing caller.
56pub struct LogSender {
57    sender: Sender<std::result::Result<(chrono::DateTime<chrono::Utc>, String), String>>,
58}
59
60impl LogSender {
61    pub fn new(
62        sender: Sender<std::result::Result<(chrono::DateTime<chrono::Utc>, String), String>>,
63    ) -> Self {
64        Self { sender }
65    }
66}
67
68#[async_trait]
69impl Logger for LogSender {
70    async fn write_log_entry(&self, msg: String) -> Result<()> {
71        info!("{msg}");
72        self.sender
73            .send(Ok((chrono::Utc::now(), msg)))
74            .await
75            .with_context(|| "LogSender::write_log_entry")
76    }
77}
78
79/// Tracing logger
80pub struct TracingLogger {}
81
82#[async_trait]
83impl Logger for TracingLogger {
84    async fn write_log_entry(&self, msg: String) -> Result<()> {
85        info!("{msg}");
86        Ok(())
87    }
88}