karo_common_rpc/
monitor.rs1use 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}