Skip to main content

krossbar_rpc/
monitor.rs

1use std::sync::atomic::{AtomicBool, Ordering};
2
3use futures::lock::Mutex;
4use log::{debug, error};
5use once_cell::sync::Lazy;
6use serde::{Deserialize, Serialize};
7use tokio::net::UnixStream;
8
9use crate::{message::RpcMessage, rpc::Rpc};
10
11static MONITOR_ACTIVE: AtomicBool = AtomicBool::new(false);
12static MONITOR_HANDLE: Lazy<Mutex<Option<Rpc>>> = Lazy::new(|| Mutex::new(None));
13
14pub const MESSAGE_METHOD: &str = "message";
15
16#[derive(Serialize, Deserialize, Debug)]
17pub enum Direction {
18    Incoming,
19    Outgoing,
20}
21
22/// Monitor message
23#[derive(Serialize, Deserialize, Debug)]
24pub struct MonitorMessage {
25    /// Peer name
26    pub peer_name: String,
27    /// Message direction
28    pub direction: Direction,
29    /// Message body
30    pub message: RpcMessage,
31}
32
33pub struct Monitor;
34
35impl Monitor {
36    pub async fn set(stream: UnixStream) {
37        debug!("Monitor connected");
38
39        *MONITOR_HANDLE.lock().await = Some(Rpc::new(stream, "monitor"));
40        MONITOR_ACTIVE.store(true, Ordering::Relaxed);
41    }
42
43    pub(crate) async fn send(message: &RpcMessage, direction: Direction, peer_name: &str) {
44        if !MONITOR_ACTIVE.load(Ordering::Relaxed) {
45            return;
46        }
47
48        let monitor_message = MonitorMessage {
49            peer_name: peer_name.to_owned(),
50            direction,
51            message: message.clone(),
52        };
53
54        if let Some(rpc) = MONITOR_HANDLE.lock().await.as_ref() {
55            if let Err(_) = rpc.send_message(MESSAGE_METHOD, &monitor_message).await {
56                MONITOR_ACTIVE.store(false, Ordering::Relaxed);
57
58                debug!("Monitor disconnected");
59            }
60        } else {
61            error!("No monitor handle when monitor indicator is active. Please submit a bug")
62        }
63    }
64}