use axum::body::Bytes;
use serde::Serialize;
use trusty_common::host_metrics::HostMetrics;
use super::service_samples::ServiceSampleBatch;
use super::transitions::ServiceTransition;
#[derive(Debug, Clone)]
pub enum HistoryEvent {
Sample(Box<HostMetrics>),
Transition(Box<ServiceTransition>),
Services(Box<ServiceSampleBatch>),
}
impl HistoryEvent {
#[must_use]
pub fn kind(&self) -> &'static str {
match self {
HistoryEvent::Sample(_) => "sample",
HistoryEvent::Transition(_) => "transition",
HistoryEvent::Services(_) => "services",
}
}
#[must_use]
pub fn frame(&self) -> Bytes {
match self {
HistoryEvent::Sample(m) => sse_frame("sample", m.as_ref()),
HistoryEvent::Transition(t) => sse_frame("transition", t.as_ref()),
HistoryEvent::Services(b) => sse_frame("services", b.as_ref()),
}
}
}
#[must_use]
pub fn sse_frame(kind: &str, payload: &impl Serialize) -> Bytes {
let json = serde_json::to_string(payload)
.unwrap_or_else(|e| format!(r#"{{"error":"serialise {kind}: {e}"}}"#));
Bytes::from(format!("event: {kind}\ndata: {json}\n\n"))
}
#[must_use]
pub fn lagged_frame(dropped: u64) -> Bytes {
sse_frame("lagged", &serde_json::json!({ "dropped": dropped }))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_frame_names_its_kind_on_one_line() {
let frame = sse_frame("transition", &serde_json::json!({ "a": 1 }));
let text = String::from_utf8(frame.to_vec()).expect("utf8 frame");
assert_eq!(text, "event: transition\ndata: {\"a\":1}\n\n");
let data_lines = text.lines().filter(|l| l.starts_with("data: ")).count();
assert_eq!(data_lines, 1, "one data line per event");
}
#[test]
fn a_sample_frame_carries_the_host_snapshot() {
let metrics = trusty_common::host_metrics::HostSampler::new().sample();
let cores = metrics.cpu.logical_cores;
let event = HistoryEvent::Sample(Box::new(metrics));
assert_eq!(event.kind(), "sample");
let text = String::from_utf8(event.frame().to_vec()).expect("utf8 frame");
let data = text
.strip_prefix("event: sample\ndata: ")
.and_then(|r| r.strip_suffix("\n\n"))
.expect("sample frame shape");
let parsed: serde_json::Value = serde_json::from_str(data).expect("sample json");
assert_eq!(parsed["cpu"]["logical_cores"], cores);
}
#[test]
fn a_services_frame_carries_the_whole_roster() {
use crate::connector::ServiceStatus;
use crate::machine_history::service_samples::{ServiceSample, ServiceSampleBatch};
let event = HistoryEvent::Services(Box::new(ServiceSampleBatch {
sampled_at_unix: 1_700_000_000,
services: vec![
ServiceSample {
id: "trusty-search".to_string(),
status: ServiceStatus::Running,
cpu_pct: Some(3.25),
rss_bytes: Some(148_897_792),
},
ServiceSample {
id: "trusty-review".to_string(),
status: ServiceStatus::Available,
cpu_pct: None,
rss_bytes: None,
},
],
}));
assert_eq!(event.kind(), "services");
let text = String::from_utf8(event.frame().to_vec()).expect("utf8 frame");
let data = text
.strip_prefix("event: services\ndata: ")
.and_then(|r| r.strip_suffix("\n\n"))
.expect("services frame shape");
let parsed: serde_json::Value = serde_json::from_str(data).expect("services json");
assert_eq!(parsed["sampled_at_unix"], 1_700_000_000_u64);
assert_eq!(parsed["services"][0]["id"], "trusty-search");
assert_eq!(parsed["services"][0]["status"], "running");
assert_eq!(parsed["services"][0]["cpu_pct"], 3.25);
assert_eq!(parsed["services"][0]["rss_bytes"], 148_897_792_u64);
assert_eq!(parsed["services"][1]["id"], "trusty-review");
assert!(
parsed["services"][1]["cpu_pct"].is_null(),
"an unmeasurable service is null, never 0.0"
);
assert!(
parsed["services"][1]["rss_bytes"].is_null(),
"an unmeasurable service is null, never 0"
);
}
#[test]
fn a_lagged_frame_carries_the_count() {
let text = String::from_utf8(lagged_frame(7).to_vec()).expect("utf8 frame");
assert_eq!(text, "event: lagged\ndata: {\"dropped\":7}\n\n");
}
}