Skip to content

Commit 02f2ea4

Browse files
committed
Introduce event streaming into a log file
1 parent 15cbc05 commit 02f2ea4

10 files changed

Lines changed: 279 additions & 10 deletions

File tree

Cargo.lock

Lines changed: 51 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

crates/hyperqueue/Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -46,6 +46,7 @@ tui = { version = "0.16" }
4646
termion = "1.5"
4747
indicatif = "0.16.2"
4848
textwrap = "0.14"
49+
async-compression = { version = "0.3.8", features = ["tokio", "gzip"] }
4950

5051
# Tako
5152
tako = { path = "../tako" }

crates/hyperqueue/src/bin/hq.rs

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -143,6 +143,10 @@ struct ServerStartOpts {
143143
/// The maximum number of events tako server will store in memory
144144
#[clap(long, default_value = "1000000")]
145145
event_store_size: usize,
146+
147+
/// Path to a log file where events will be stored.
148+
#[clap(long, hide(true))]
149+
event_log_file: Option<PathBuf>,
146150
}
147151

148152
#[derive(Parser)]
@@ -292,7 +296,8 @@ async fn command_server_start(
292296
autoalloc_interval: opts.autoalloc_interval.map(|x| x.unpack()),
293297
client_port: opts.client_port,
294298
worker_port: opts.worker_port,
295-
event_store_size: opts.event_store_size,
299+
event_buffer_size: opts.event_store_size,
300+
event_log_file: opts.event_log_file,
296301
};
297302

298303
init_hq_server(gsettings, server_cfg).await
Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,24 @@
1+
mod stream;
2+
mod write;
3+
4+
pub use stream::{start_event_streaming, EventStreamSender};
5+
pub use write::EventLogWriter;
6+
7+
use bstr::BString;
8+
use serde::{Deserialize, Serialize};
9+
10+
const HQ_LOG_HEADER: &[u8] = b"hq-event-log";
11+
const HQ_LOG_VERSION: u32 = 0;
12+
13+
fn canonical_header() -> LogFileHeader {
14+
LogFileHeader {
15+
header: HQ_LOG_HEADER.into(),
16+
version: HQ_LOG_VERSION,
17+
}
18+
}
19+
20+
#[derive(Serialize, Deserialize, Debug, Eq, PartialEq)]
21+
struct LogFileHeader {
22+
header: BString,
23+
version: u32,
24+
}
Lines changed: 70 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,70 @@
1+
use crate::event::log::write::EventLogWriter;
2+
use crate::event::MonitoringEvent;
3+
use std::future::Future;
4+
use std::time::Duration;
5+
use tokio::sync::mpsc;
6+
7+
pub type EventStreamSender = mpsc::UnboundedSender<MonitoringEvent>;
8+
pub type EventStreamReceiver = mpsc::UnboundedReceiver<MonitoringEvent>;
9+
10+
fn create_event_stream_queue() -> (EventStreamSender, EventStreamReceiver) {
11+
mpsc::unbounded_channel()
12+
}
13+
14+
/// Start event streaming into a log file.
15+
/// Streaming is running on another thread to reduce overhead and interference.
16+
///
17+
/// Returns a future that resolves once the event streaming thread finishes.
18+
/// The thread will finish if there is some I/O error or if the `receiver` is closed.
19+
pub fn start_event_streaming(
20+
writer: EventLogWriter,
21+
) -> (EventStreamSender, impl Future<Output = ()>) {
22+
let (tx, rx) = create_event_stream_queue();
23+
24+
let handle = std::thread::spawn(move || {
25+
let process = streaming_process(writer, rx);
26+
27+
let runtime = tokio::runtime::Builder::new_current_thread()
28+
.enable_all()
29+
.build()
30+
.unwrap();
31+
32+
if let Err(error) = runtime.block_on(process) {
33+
log::error!("Event streaming has ended with an error: {error:?}");
34+
} else {
35+
log::debug!("Event streaming has finished successfully");
36+
}
37+
});
38+
let end_fut = async move {
39+
handle.join().expect("Event streaming thread has crashed");
40+
};
41+
(tx, end_fut)
42+
}
43+
44+
const FLUSH_PERIOD: Duration = Duration::from_secs(30);
45+
46+
async fn streaming_process(
47+
mut writer: EventLogWriter,
48+
mut receiver: EventStreamReceiver,
49+
) -> anyhow::Result<()> {
50+
let mut flush_fut = tokio::time::interval(FLUSH_PERIOD);
51+
52+
loop {
53+
tokio::select! {
54+
_ = flush_fut.tick() => {
55+
writer.flush().await?;
56+
}
57+
res = receiver.recv() => {
58+
match res {
59+
Some(event) => {
60+
log::trace!("Event: {event:?}");
61+
writer.store(event).await?;
62+
}
63+
None => break
64+
}
65+
}
66+
}
67+
}
68+
writer.finish().await?;
69+
Ok(())
70+
}
Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
1+
use crate::event::log::canonical_header;
2+
use crate::event::MonitoringEvent;
3+
use async_compression::tokio::write::GzipEncoder;
4+
use async_compression::Level;
5+
use std::path::Path;
6+
use tokio::fs::File;
7+
use tokio::io::AsyncWriteExt;
8+
9+
/// Streams monitoring events into a file on disk.
10+
pub struct EventLogWriter {
11+
file: GzipEncoder<File>,
12+
buffer: Vec<u8>,
13+
}
14+
15+
const BUF_MAX_SIZE: usize = 16 * 1024;
16+
17+
impl EventLogWriter {
18+
pub async fn create(path: &Path) -> anyhow::Result<Self> {
19+
let mut file = File::create(path).await?;
20+
let header = rmp_serde::encode::to_vec(&canonical_header())?;
21+
file.write_all(&header).await?;
22+
file.flush().await?;
23+
24+
let file = GzipEncoder::with_quality(file, Level::Fastest);
25+
26+
// Keep buffer capacity larger than max size to avoid reallocation if we overflow
27+
// the buffer.
28+
let buffer = Vec::with_capacity(BUF_MAX_SIZE * 2);
29+
Ok(Self { file, buffer })
30+
}
31+
32+
#[inline]
33+
pub async fn store(&mut self, event: MonitoringEvent) -> anyhow::Result<()> {
34+
rmp_serde::encode::write(&mut self.buffer, &event)?;
35+
if self.is_buffer_full() {
36+
self.write_buffer().await?;
37+
}
38+
Ok(())
39+
}
40+
41+
pub async fn flush(&mut self) -> anyhow::Result<()> {
42+
if !self.buffer.is_empty() {
43+
self.write_buffer().await?;
44+
}
45+
self.file.flush().await?;
46+
Ok(())
47+
}
48+
49+
pub async fn finish(mut self) -> anyhow::Result<()> {
50+
self.flush().await?;
51+
self.file.shutdown().await?;
52+
Ok(())
53+
}
54+
55+
async fn write_buffer(&mut self) -> tokio::io::Result<()> {
56+
self.file.write_all(&self.buffer).await?;
57+
self.buffer.clear();
58+
Ok(())
59+
}
60+
61+
#[inline]
62+
fn is_buffer_full(&self) -> bool {
63+
self.buffer.len() >= BUF_MAX_SIZE
64+
}
65+
}

crates/hyperqueue/src/event/mod.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
pub mod events;
2+
pub mod log;
23
pub mod storage;
34

45
use events::MonitoringEventPayload;

crates/hyperqueue/src/event/storage.rs

Lines changed: 20 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
use crate::event::events::MonitoringEventPayload;
2+
use crate::event::log::EventStreamSender;
23
use crate::event::{MonitoringEvent, MonitoringEventId};
34
use crate::WorkerId;
45
use std::collections::VecDeque;
@@ -11,6 +12,7 @@ pub struct EventStorage {
1112
event_store_size: usize,
1213
event_queue: VecDeque<MonitoringEvent>,
1314
last_event_id: u32,
15+
stream_sender: Option<EventStreamSender>,
1416
}
1517

1618
impl Default for EventStorage {
@@ -19,16 +21,18 @@ impl Default for EventStorage {
1921
event_store_size: 1_000_000,
2022
event_queue: Default::default(),
2123
last_event_id: 0,
24+
stream_sender: None,
2225
}
2326
}
2427
}
2528

2629
impl EventStorage {
27-
pub fn new(event_store_size: usize) -> Self {
30+
pub fn new(event_store_size: usize, stream_sender: Option<EventStreamSender>) -> Self {
2831
Self {
2932
event_store_size,
3033
event_queue: VecDeque::new(),
3134
last_event_id: 0,
35+
stream_sender,
3236
}
3337
}
3438

@@ -65,13 +69,26 @@ impl EventStorage {
6569

6670
fn insert_event(&mut self, payload: MonitoringEventPayload) {
6771
self.last_event_id += 1;
68-
self.event_queue.push_back(MonitoringEvent {
72+
73+
let event = MonitoringEvent {
6974
payload,
7075
id: self.last_event_id,
7176
time: SystemTime::now(),
72-
});
77+
};
78+
self.stream_event(&event);
79+
self.event_queue.push_back(event);
80+
7381
if self.event_queue.len() > self.event_store_size {
7482
self.event_queue.pop_front();
7583
}
7684
}
85+
86+
fn stream_event(&mut self, event: &MonitoringEvent) {
87+
if let Some(ref streamer) = self.stream_sender {
88+
if streamer.send(event.clone()).is_err() {
89+
log::error!("Event streaming queue has been closed.");
90+
self.stream_sender = None;
91+
}
92+
}
93+
}
7794
}

crates/hyperqueue/src/lib.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ pub type Result<T> = std::result::Result<T, Error>;
2020
// ID types
2121
use tako::define_id_type;
2222

23-
pub type WorkerId = tako::WorkerId;
23+
pub use tako::WorkerId;
2424
pub type TakoTaskId = tako::TaskId;
2525
pub type Priority = tako::Priority;
2626

0 commit comments

Comments
 (0)