From e02efe9ba06820e2a4f4e7e9f5f10c05ef4b8442 Mon Sep 17 00:00:00 2001 From: Omer Faruk Bayram Date: Sun, 13 Aug 2023 15:13:59 +0000 Subject: [PATCH] event_monitor: make it possible to `subscribe` to `Monitor` Signed-off-by: Omer Faruk Bayram --- event_monitor/src/lib.rs | 20 +++++++++++++++++++- vmm/src/lib.rs | 6 ++++++ 2 files changed, 25 insertions(+), 1 deletion(-) diff --git a/event_monitor/src/lib.rs b/event_monitor/src/lib.rs index 68f878395..b8f9b9ab9 100644 --- a/event_monitor/src/lib.rs +++ b/event_monitor/src/lib.rs @@ -9,6 +9,7 @@ use std::collections::HashMap; use std::fs::File; use std::io; use std::os::unix::io::AsRawFd; +use std::sync::Arc; use std::time::{Duration, Instant}; static mut MONITOR: Option = None; @@ -24,6 +25,23 @@ struct Event<'a> { pub struct Monitor { pub rx: flume::Receiver, pub file: File, + pub broadcast: Vec>>, +} + +impl Monitor { + pub fn new(rx: flume::Receiver, file: File) -> Self { + Self { + rx, + file, + broadcast: vec![], + } + } + + pub fn subscribe(&mut self) -> flume::Receiver> { + let (tx, rx) = flume::unbounded(); + self.broadcast.push(tx); + rx + } } struct MonitorHandle { @@ -56,7 +74,7 @@ pub fn set_monitor(file: File) -> io::Result { set_file_nonblocking(&file)?; let (tx, rx) = flume::unbounded(); - let monitor = Monitor { rx, file }; + let monitor = Monitor::new(rx, file); // SAFETY: MONITOR is None. Nobody else can hold a reference to it. unsafe { diff --git a/vmm/src/lib.rs b/vmm/src/lib.rs index 7cad2ca42..06d13cd51 100644 --- a/vmm/src/lib.rs +++ b/vmm/src/lib.rs @@ -319,8 +319,14 @@ pub fn start_event_monitor_thread( std::panic::catch_unwind(AssertUnwindSafe(move || { while let Ok(event) = monitor.rx.recv() { + let event = Arc::new(event); + monitor.file.write_all(event.as_bytes().as_ref()).ok(); monitor.file.write_all(b"\n\n").ok(); + + for tx in monitor.broadcast.iter() { + tx.send(event.clone()).ok(); + } } })) .map_err(|_| {