Skip to main content

karo_common_rpc/
monitor.rs

1use std::sync::atomic::{AtomicBool, Ordering};
2
3use futures::{lock::Mutex, Stream};
4use log::{debug, error};
5use once_cell::sync::Lazy;
6use serde::{Deserialize, Serialize};
7use tokio::net::UnixStream;
8
9use crate::{message::RpcMessage, message_stream::AsyncWriteMessage};
10
11static MONITOR_ACTIVE: AtomicBool = AtomicBool::new(false);
12static MONITOR_HANDLE: Lazy<Mutex<Option<UnixStream>>> = Lazy::new(|| Mutex::new(None));
13
14#[derive(Serialize, Deserialize, Debug)]
15pub enum Direction {
16    Incoming,
17    Ougoing,
18}
19
20#[derive(Serialize, Deserialize, Debug)]
21pub struct MonitorMessage {
22    pub direction: Direction,
23    pub message: RpcMessage,
24}
25
26pub struct Monitor;
27
28impl Monitor {
29    pub async fn set(stream: UnixStream) {
30        debug!("Monitor connected");
31
32        *MONITOR_HANDLE.lock().await = Some(stream);
33        MONITOR_ACTIVE.store(true, Ordering::Relaxed);
34    }
35
36    #[cfg(feature = "impl-monitor")]
37    pub fn make_receiver(mut stream: UnixStream) -> impl Stream<Item = MonitorMessage> {
38        use std::pin::pin;
39
40        use futures::{stream, Future};
41
42        use crate::message_stream::AsyncReadMessage;
43
44        stream::poll_fn(move |ctx| {
45            pin!(stream.read_message()).poll(ctx).map(|val| match val {
46                Ok(message) => Some(message),
47                Err(_) => None,
48            })
49        })
50    }
51
52    pub(crate) async fn send(message: &RpcMessage, direction: Direction) {
53        if !MONITOR_ACTIVE.load(Ordering::Relaxed) {
54            return;
55        }
56
57        let monitor_message = MonitorMessage {
58            direction,
59            message: message.clone(),
60        };
61
62        let mut lock = MONITOR_HANDLE.lock().await;
63        if let Some(stream) = lock.as_mut() {
64            if let Err(_) = stream.write_message(&monitor_message).await {
65                MONITOR_ACTIVE.store(false, Ordering::Relaxed);
66
67                debug!("Monitor disconnected");
68            }
69        } else {
70            error!("No monitor handle when monitor indicator is active. Please submit a bug")
71        }
72    }
73}