micromegas_analytics/
response_writer.rs1use 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#[async_trait]
10pub trait Logger: Send + Sync {
11 async fn write_log_entry(&self, msg: String) -> Result<()>;
12}
13
14pub 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
51pub 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
79pub 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}