pub mod events;
pub mod sampler;
pub mod service_samples;
pub mod stream;
pub mod transitions;
use std::collections::BTreeMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use serde::{Deserialize, Serialize};
use tokio::sync::{RwLock, broadcast};
use trusty_common::console_metrics::ConsoleMetricsReport;
use trusty_common::host_metrics::HostMetrics;
use trusty_common::host_metrics::history::{
HOST_HISTORY_CAPACITY, HOST_SAMPLE_INTERVAL_SECS, MetricRing,
};
use events::HistoryEvent;
use service_samples::{ServiceSample, ServiceSampleBatch};
use transitions::{SERVICE_REPORT_GRACE_SECS, ServiceTransition, TransitionTracker};
pub const MACHINE_HISTORY_SCHEMA_VERSION: u32 = 4;
pub const TRANSITION_LOG_CAPACITY: usize = 256;
pub const EVENT_BUFFER: usize = HOST_HISTORY_CAPACITY * 2;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HistorySnapshot {
pub samples: Vec<HostMetrics>,
pub transitions: Vec<ServiceTransition>,
pub service_samples: BTreeMap<String, Vec<ServiceSample>>,
pub sample_capacity: usize,
pub service_sample_capacity: usize,
pub transition_capacity: usize,
pub sample_interval_secs: u64,
pub sse_client_count: usize,
pub schema_version: u32,
}
struct Inner {
samples: MetricRing<HostMetrics>,
transitions: MetricRing<ServiceTransition>,
service_samples: BTreeMap<String, MetricRing<ServiceSample>>,
service_capacity: usize,
tracker: TransitionTracker,
}
struct Shared {
inner: RwLock<Inner>,
events: broadcast::Sender<HistoryEvent>,
sample_interval_secs: AtomicU64,
}
#[derive(Clone)]
pub struct MachineHistory {
shared: Arc<Shared>,
}
impl Default for MachineHistory {
fn default() -> Self {
Self::new()
}
}
impl MachineHistory {
#[must_use]
pub fn new() -> Self {
Self::with_limits(
HOST_HISTORY_CAPACITY,
TRANSITION_LOG_CAPACITY,
EVENT_BUFFER,
Duration::from_secs(SERVICE_REPORT_GRACE_SECS),
)
}
#[must_use]
pub fn with_limits(
sample_capacity: usize,
transition_capacity: usize,
event_buffer: usize,
grace: Duration,
) -> Self {
let (events, _) = broadcast::channel(event_buffer.max(1));
Self {
shared: Arc::new(Shared {
inner: RwLock::new(Inner {
samples: MetricRing::new(sample_capacity),
transitions: MetricRing::new(transition_capacity),
service_samples: BTreeMap::new(),
service_capacity: sample_capacity,
tracker: TransitionTracker::new(grace),
}),
events,
sample_interval_secs: AtomicU64::new(HOST_SAMPLE_INTERVAL_SECS),
}),
}
}
pub fn set_sample_interval(&self, secs: u64) {
self.shared
.sample_interval_secs
.store(secs, Ordering::Relaxed);
}
pub async fn record_sample(&self, sample: HostMetrics) {
let mut inner = self.shared.inner.write().await;
inner.samples.push(sample.clone());
let _ = self
.shared
.events
.send(HistoryEvent::Sample(Box::new(sample)));
}
pub async fn record_service_samples(&self, batch: ServiceSampleBatch) {
let mut inner = self.shared.inner.write().await;
let capacity = inner.service_capacity;
for sample in &batch.services {
inner
.service_samples
.entry(sample.id.clone())
.or_insert_with(|| MetricRing::new(capacity))
.push(sample.clone());
}
let _ = self
.shared
.events
.send(HistoryEvent::Services(Box::new(batch)));
}
pub async fn latest_service_metrics(&self) -> std::collections::HashMap<String, ServiceSample> {
let inner = self.shared.inner.read().await;
inner
.service_samples
.iter()
.filter_map(|(id, ring)| ring.last().map(|s| (id.clone(), s.clone())))
.collect()
}
pub async fn observe_services(
&self,
reports: &[ConsoleMetricsReport],
) -> Vec<ServiceTransition> {
let now_unix = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_secs());
let mut inner = self.shared.inner.write().await;
let changes = inner.tracker.observe(reports, Instant::now(), now_unix);
for change in &changes {
inner.transitions.push(change.clone());
let _ = self
.shared
.events
.send(HistoryEvent::Transition(Box::new(change.clone())));
}
changes
}
pub async fn snapshot(&self) -> HistorySnapshot {
let inner = self.shared.inner.read().await;
self.snapshot_locked(&inner)
}
pub async fn subscribe(&self) -> (HistorySnapshot, broadcast::Receiver<HistoryEvent>) {
let inner = self.shared.inner.read().await;
let rx = self.shared.events.subscribe();
(self.snapshot_locked(&inner), rx)
}
#[must_use]
pub fn subscriber_count(&self) -> usize {
self.shared.events.receiver_count()
}
fn snapshot_locked(&self, inner: &Inner) -> HistorySnapshot {
HistorySnapshot {
samples: inner.samples.snapshot(),
transitions: inner.transitions.snapshot(),
service_samples: inner
.service_samples
.iter()
.map(|(id, ring)| (id.clone(), ring.snapshot()))
.collect(),
sample_capacity: inner.samples.capacity(),
service_sample_capacity: inner.service_capacity,
transition_capacity: inner.transitions.capacity(),
sample_interval_secs: self.shared.sample_interval_secs.load(Ordering::Relaxed),
sse_client_count: self.subscriber_count(),
schema_version: MACHINE_HISTORY_SCHEMA_VERSION,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use trusty_common::console_metrics::{ServiceHealth, make_report};
use trusty_common::host_metrics::HostSampler;
fn sample() -> HostMetrics {
HostSampler::new().sample()
}
#[tokio::test]
async fn history_starts_empty() {
let h = MachineHistory::new();
let snap = h.snapshot().await;
assert!(snap.samples.is_empty());
assert!(snap.transitions.is_empty());
assert_eq!(snap.sample_capacity, HOST_HISTORY_CAPACITY);
assert_eq!(snap.transition_capacity, TRANSITION_LOG_CAPACITY);
assert_eq!(snap.sample_interval_secs, HOST_SAMPLE_INTERVAL_SECS);
assert_eq!(snap.schema_version, MACHINE_HISTORY_SCHEMA_VERSION);
assert_eq!(snap.sse_client_count, 0);
}
#[tokio::test]
async fn subscriber_count_moves_by_one_per_stream() {
let h = MachineHistory::new();
let baseline = h.subscriber_count();
assert_eq!(baseline, 0, "a fresh history has no streams attached");
let (snap, first) = h.subscribe().await;
assert_eq!(h.subscriber_count(), baseline + 1);
assert_eq!(
snap.sse_client_count,
h.subscriber_count(),
"the snapshot field must be the live count, not a second tally"
);
let (_, second) = h.subscribe().await;
assert_eq!(
h.subscriber_count(),
baseline + 2,
"a second stream raises the count by exactly one"
);
drop(second);
assert_eq!(
h.subscriber_count(),
baseline + 1,
"closing one stream lowers the count by exactly one"
);
drop(first);
assert_eq!(
h.subscriber_count(),
baseline,
"closing the last stream returns the count to the baseline"
);
assert_eq!(h.snapshot().await.sse_client_count, baseline);
}
#[tokio::test]
async fn recording_a_sample_fans_out_to_subscribers() {
let h = MachineHistory::new();
let (snap, mut rx) = h.subscribe().await;
assert!(snap.samples.is_empty());
h.record_sample(sample()).await;
match rx.try_recv() {
Ok(HistoryEvent::Sample(_)) => {}
other => panic!("expected a sample event, got {other:?}"),
}
assert_eq!(h.snapshot().await.samples.len(), 1);
}
#[tokio::test]
async fn the_ring_bounds_what_history_returns() {
let h = MachineHistory::with_limits(3, 4, 8, Duration::from_secs(60));
for _ in 0..4 {
h.record_sample(sample()).await;
}
let snap = h.snapshot().await;
assert_eq!(snap.samples.len(), 3, "the ring bounds the window");
assert_eq!(snap.sample_capacity, 3);
assert_eq!(snap.transition_capacity, 4);
}
#[tokio::test]
async fn an_unchanged_service_adds_nothing_to_the_log() {
let h = MachineHistory::new();
let ok = make_report(
"trusty-search",
"Trusty Search",
"1.0.0",
ServiceHealth::Ok,
serde_json::json!({}),
1,
);
for _ in 0..5 {
assert!(
h.observe_services(std::slice::from_ref(&ok))
.await
.is_empty()
);
}
assert!(h.snapshot().await.transitions.is_empty());
let degraded = make_report(
"trusty-search",
"Trusty Search",
"1.0.0",
ServiceHealth::Degraded,
serde_json::json!({}),
1,
);
let changes = h.observe_services(&[degraded]).await;
assert_eq!(changes.len(), 1);
let snap = h.snapshot().await;
assert_eq!(snap.transitions.len(), 1);
assert_eq!(snap.transitions[0].to, transitions::ServiceState::Degraded);
}
fn service_batch(id: &str, cpu: Option<f32>, at: u64) -> ServiceSampleBatch {
use crate::connector::ServiceStatus;
ServiceSampleBatch {
sampled_at_unix: at,
services: vec![ServiceSample {
id: id.to_string(),
status: ServiceStatus::Running,
cpu_pct: cpu,
rss_bytes: cpu.map(|c| (c as u64 + 1) * 1024 * 1024),
}],
}
}
#[tokio::test]
async fn recording_service_samples_fans_out_to_subscribers() {
let h = MachineHistory::new();
let (snap, mut rx) = h.subscribe().await;
assert!(snap.service_samples.is_empty());
h.record_service_samples(service_batch("trusty-search", Some(2.0), 10))
.await;
match rx.try_recv() {
Ok(HistoryEvent::Services(b)) => {
assert_eq!(b.sampled_at_unix, 10);
assert_eq!(b.services.len(), 1);
}
other => panic!("expected a services event, got {other:?}"),
}
let snap = h.snapshot().await;
assert_eq!(snap.service_samples["trusty-search"].len(), 1);
}
#[tokio::test]
async fn the_service_ring_bounds_what_history_returns() {
let h = MachineHistory::with_limits(3, 4, 8, Duration::from_secs(60));
for at in 1..=4 {
h.record_service_samples(service_batch("trusty-search", Some(at as f32), at))
.await;
}
let snap = h.snapshot().await;
let ring = &snap.service_samples["trusty-search"];
assert_eq!(ring.len(), 3, "the ring bounds the per-service window");
assert_eq!(
ring.first().and_then(|s| s.cpu_pct),
Some(2.0),
"the oldest sample was evicted"
);
assert_eq!(snap.service_sample_capacity, 3);
}
#[tokio::test]
async fn the_snapshot_carries_a_ring_per_service() {
let h = MachineHistory::new();
h.record_service_samples(service_batch("trusty-search", Some(1.0), 1))
.await;
h.record_service_samples(service_batch("trusty-search", Some(2.0), 2))
.await;
h.record_service_samples(service_batch("trusty-mpm", None, 3))
.await;
let snap = h.snapshot().await;
assert_eq!(snap.service_samples.len(), 2);
assert_eq!(snap.service_samples["trusty-search"].len(), 2);
assert_eq!(snap.service_samples["trusty-mpm"].len(), 1);
assert_eq!(snap.schema_version, 4, "#6908 bumped the payload shape");
}
#[tokio::test]
async fn latest_service_metrics_reads_the_newest_sample() {
let h = MachineHistory::new();
h.record_service_samples(service_batch("trusty-search", Some(1.0), 1))
.await;
h.record_service_samples(service_batch("trusty-search", Some(9.5), 2))
.await;
h.record_service_samples(service_batch("trusty-mpm", Some(3.0), 2))
.await;
let latest = h.latest_service_metrics().await;
assert_eq!(latest["trusty-search"].cpu_pct, Some(9.5));
assert_eq!(latest["trusty-search"].rss_bytes, Some(10 * 1024 * 1024));
h.record_service_samples(service_batch("trusty-search", None, 3))
.await;
let latest = h.latest_service_metrics().await;
assert_eq!(latest["trusty-search"].cpu_pct, None);
assert_eq!(latest["trusty-search"].rss_bytes, None);
assert_eq!(latest["trusty-mpm"].cpu_pct, Some(3.0));
}
#[tokio::test]
async fn a_subscriber_stalled_for_the_whole_window_is_not_lagged() {
use tokio::sync::broadcast::error::TryRecvError;
const TICKS: u64 = 200;
let h = MachineHistory::new();
let (_snap, mut rx) = h.subscribe().await;
for seq in 0..TICKS {
h.record_sample(tagged_sample(seq)).await;
h.record_service_samples(service_batch("trusty-search", Some(1.0), seq))
.await;
}
let mut samples = 0u64;
let mut batches = 0u64;
loop {
match rx.try_recv() {
Ok(HistoryEvent::Sample(m)) => {
assert_eq!(
m.sampled_at_unix,
Some(samples),
"samples must arrive in order with no gap"
);
samples += 1;
}
Ok(HistoryEvent::Services(b)) => {
assert_eq!(b.sampled_at_unix, batches);
batches += 1;
}
Ok(HistoryEvent::Transition(_)) => {}
Err(TryRecvError::Empty | TryRecvError::Closed) => break,
Err(TryRecvError::Lagged(n)) => panic!(
"a subscriber stalled for {TICKS} ticks ({} events) was told it lagged by \
{n}; EVENT_BUFFER is {EVENT_BUFFER}, which must cover two events per tick \
across the whole {HOST_HISTORY_CAPACITY}-point window",
TICKS * 2
),
}
}
assert_eq!(samples, TICKS, "every host sample survived the stall");
assert_eq!(batches, TICKS, "every service batch survived the stall");
}
#[tokio::test]
async fn the_advertised_interval_follows_the_sampler() {
let h = MachineHistory::new();
h.set_sample_interval(30);
assert_eq!(h.snapshot().await.sample_interval_secs, 30);
}
fn tagged_sample(seq: u64) -> HostMetrics {
use trusty_common::host_metrics::{
CpuMetrics, DiskMetrics, MemoryMetrics, NetworkMetrics, Pressure,
};
HostMetrics {
cpu: CpuMetrics {
usage_pct: 0.0,
logical_cores: 1,
physical_cores: None,
pressure: Pressure::Nominal,
},
memory: MemoryMetrics {
total_bytes: 1,
used_bytes: 0,
available_bytes: 1,
usage_pct: 0.0,
swap_total_bytes: 0,
swap_used_bytes: 0,
pressure: Pressure::Nominal,
},
disks: DiskMetrics {
aggregate_total_bytes: 1,
aggregate_available_bytes: 1,
aggregate_used_bytes: 0,
aggregate_usage_pct: 0.0,
pressure: Pressure::Nominal,
mounts: Vec::new(),
},
network: NetworkMetrics {
rx_bytes_per_sec: 0.0,
tx_bytes_per_sec: 0.0,
rx_total_bytes: 0,
tx_total_bytes: 0,
window_secs: 1.0,
},
overall_pressure: Pressure::Nominal,
sampled_at_unix: Some(seq),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn every_sample_reaches_a_mid_run_subscriber_exactly_once() {
use tokio::sync::broadcast::error::TryRecvError;
const TOTAL: u64 = 20_000;
const MAX_SUBSCRIPTIONS: usize = 200_000;
let history = MachineHistory::with_limits(
8,
8,
TOTAL as usize * 2,
Duration::from_secs(60),
);
let writer = {
let history = history.clone();
tokio::spawn(async move {
for seq in 0..TOTAL {
history.record_sample(tagged_sample(seq)).await;
}
})
};
let mut taken: Vec<(Option<u64>, broadcast::Receiver<HistoryEvent>)> = Vec::new();
while !writer.is_finished() && taken.len() < MAX_SUBSCRIPTIONS {
let (snapshot, rx) = history.subscribe().await;
let last = snapshot
.samples
.last()
.map(|m| m.sampled_at_unix.expect("tagged sample"));
taken.push((last, rx));
}
writer.await.expect("writer task");
let mut mid_run = 0usize;
for (nth, (last_snapshotted, mut rx)) in taken.into_iter().enumerate() {
if last_snapshotted == Some(TOTAL - 1) {
continue;
}
let expected = last_snapshotted.map_or(0, |s| s + 1);
mid_run += 1;
let first_live = loop {
match rx.try_recv() {
Ok(HistoryEvent::Sample(m)) => break m.sampled_at_unix.expect("tagged sample"),
Ok(HistoryEvent::Transition(_) | HistoryEvent::Services(_)) => {}
Err(TryRecvError::Empty | TryRecvError::Closed) => panic!(
"subscription {nth} snapshotted through {last_snapshotted:?} of {TOTAL} \
samples and then received nothing live"
),
Err(TryRecvError::Lagged(n)) => panic!(
"subscription {nth} lagged by {n} — the buffer was sized to prevent it"
),
}
};
assert_eq!(
first_live, expected,
"subscription {nth} snapshotted through {last_snapshotted:?} and must resume live \
at {expected}; resuming later means a sample fell between the snapshot and the \
subscription, resuming at the same sequence means it landed in both"
);
}
assert!(
mid_run >= 500,
"only {mid_run} subscription(s) landed mid-run — the race this test exists to catch \
was never given a chance"
);
}
}