use std::sync::Arc;
use tokio::sync::{Mutex, RwLock, mpsc};
use tokio::time::{Duration, interval};
use crate::correlate::Trace;
use crate::correlate::window::TraceWindow;
use crate::detect;
#[cfg(test)]
use crate::detect::sanitizer_aware::SanitizerAwareMode;
use crate::detect::{Confidence, DetectConfig};
use crate::event::SpanEvent;
use crate::normalize;
use crate::report::metrics::MetricsState;
use crate::report::{DatabaseWaste, GreenSummary, MessagingWaste};
use crate::score;
use crate::score::alumet::{AlumetState, DbEnergyState};
use crate::score::cloud_energy::CloudEnergyState;
use crate::score::electricity_maps::ElectricityMapsState;
use crate::score::kepler::KeplerState;
use crate::score::redfish::RedfishState;
use crate::score::scaphandre::ScaphandreState;
use super::findings_store;
use super::sampling::apply_sampling;
#[derive(Clone, Copy)]
pub(super) struct EventLoopConfig {
pub(super) green_enabled: bool,
pub(super) sampling_rate: f64,
pub(super) evict_ms: u64,
pub(super) confidence: Confidence,
pub(super) waste_sticky_ttl_ms: u64,
pub(super) analysis_queue_capacity: usize,
}
pub(super) struct ShutdownTargets<'a> {
pub(super) energy: EnergyScraperHandles<'a>,
pub(super) listeners: ListenerHandles<'a>,
}
#[derive(Clone, Copy)]
pub(super) struct EnergyScraperHandles<'a> {
pub(super) alumet: Option<&'a tokio::task::JoinHandle<()>>,
pub(super) scaphandre: Option<&'a tokio::task::JoinHandle<()>>,
pub(super) kepler: Option<&'a tokio::task::JoinHandle<()>>,
pub(super) redfish: Option<&'a tokio::task::JoinHandle<()>>,
pub(super) cloud: Option<&'a tokio::task::JoinHandle<()>>,
pub(super) emaps: Option<&'a tokio::task::JoinHandle<()>>,
}
#[derive(Clone, Copy)]
pub(super) struct ListenerHandles<'a> {
pub(super) grpc: &'a tokio::task::JoinHandle<()>,
pub(super) http: &'a tokio::task::JoinHandle<()>,
pub(super) json_socket: Option<&'a tokio::task::JoinHandle<()>>,
}
pub(super) struct EnergySources<'a> {
pub(super) base_carbon_ctx: Arc<score::carbon::CarbonContext>,
pub(super) alumet_state: Option<&'a AlumetState>,
pub(super) alumet_db_state: Option<&'a DbEnergyState>,
pub(super) alumet_broker_state: Option<&'a DbEnergyState>,
pub(super) static_broker: Option<(
&'a score::broker_static::StaticBrokerConfig,
&'a score::broker_static::StaticBrokerState,
)>,
pub(super) alumet_staleness_ms: u64,
pub(super) scaphandre_state: Option<&'a ScaphandreState>,
pub(super) scaphandre_staleness_ms: u64,
pub(super) kepler_state: Option<&'a KeplerState>,
pub(super) kepler_staleness_ms: u64,
pub(super) redfish_state: Option<&'a RedfishState>,
pub(super) redfish_staleness_ms: u64,
pub(super) cloud_state: Option<&'a CloudEnergyState>,
pub(super) cloud_staleness_ms: u64,
pub(super) emaps_state: Option<&'a ElectricityMapsState>,
pub(super) emaps_staleness_ms: u64,
}
struct AnalysisBatch {
traces: Vec<(String, Vec<normalize::NormalizedEvent>)>,
carbon_ctx: Arc<score::carbon::CarbonContext>,
}
impl AnalysisBatch {
fn new(
traces: Vec<(String, Vec<normalize::NormalizedEvent>)>,
sources: &EnergySources<'_>,
) -> Self {
Self {
traces,
carbon_ctx: build_owned_tick_ctx(sources),
}
}
}
struct AnalysisWorkerCtx {
detect_config: DetectConfig,
green_enabled: bool,
confidence: Confidence,
metrics: Arc<MetricsState>,
findings_store: Arc<findings_store::FindingsStore>,
traces_store: Arc<super::traces_store::TracesStore>,
correlator: Option<Arc<Mutex<detect::correlate_cross::CrossTraceCorrelator>>>,
green_summary_cell: Arc<RwLock<GreenSummary>>,
archive_tx: Option<mpsc::Sender<super::archive::OwnedArchive>>,
waste_sticky_ttl_ms: u64,
}
#[allow(clippy::too_many_arguments)]
pub(super) async fn run_event_loop(
rx: &mut mpsc::Receiver<Vec<SpanEvent>>,
window: &Arc<Mutex<TraceWindow>>,
metrics: Arc<MetricsState>,
findings_store: Arc<findings_store::FindingsStore>,
traces_store: Arc<super::traces_store::TracesStore>,
correlator: Option<Arc<Mutex<detect::correlate_cross::CrossTraceCorrelator>>>,
detect_config: &DetectConfig,
energy_sources: &EnergySources<'_>,
shutdown: ShutdownTargets<'_>,
loop_cfg: EventLoopConfig,
green_summary_cell: Arc<RwLock<GreenSummary>>,
archive_tx: Option<mpsc::Sender<super::archive::OwnedArchive>>,
) -> Result<(), super::DaemonError> {
let (work_tx, work_rx) = mpsc::channel::<AnalysisBatch>(loop_cfg.analysis_queue_capacity);
let worker = tokio::spawn(run_analysis_worker(
work_rx,
AnalysisWorkerCtx {
detect_config: detect_config.clone(),
green_enabled: loop_cfg.green_enabled,
confidence: loop_cfg.confidence,
metrics: metrics.clone(),
findings_store,
traces_store,
correlator,
green_summary_cell,
archive_tx,
waste_sticky_ttl_ms: loop_cfg.waste_sticky_ttl_ms,
},
));
drive_event_loop(
rx,
window,
&metrics,
energy_sources,
shutdown,
loop_cfg,
work_tx,
worker,
crate::shutdown::shutdown_signal(),
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn drive_event_loop(
rx: &mut mpsc::Receiver<Vec<SpanEvent>>,
window: &Arc<Mutex<TraceWindow>>,
metrics: &MetricsState,
energy_sources: &EnergySources<'_>,
shutdown: ShutdownTargets<'_>,
loop_cfg: EventLoopConfig,
work_tx: mpsc::Sender<AnalysisBatch>,
mut worker: tokio::task::JoinHandle<()>,
shutdown_fut: impl Future<Output = ()>,
) -> Result<(), super::DaemonError> {
let mut ticker = interval(Duration::from_millis(loop_cfg.evict_ms.max(100)));
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let mut service_meter = ServiceMeter {
known_services: std::collections::HashMap::new(),
max_service_cardinality: MAX_SERVICE_CARDINALITY,
service_cap_warned: false,
};
tokio::pin!(shutdown_fut);
let graceful = loop {
tokio::select! {
Some(events) = rx.recv() => {
let lru_evicted = ingest_event_batch(
events,
loop_cfg.sampling_rate,
window,
metrics,
&mut service_meter,
).await;
enqueue_for_analysis(lru_evicted, energy_sources, &work_tx, metrics);
}
_ = ticker.tick() => {
let expired = evict_expired_traces(window, metrics).await;
enqueue_for_analysis(expired, energy_sources, &work_tx, metrics);
}
() = &mut shutdown_fut => {
tracing::info!("Shutting down daemon, processing remaining traces...");
break true;
}
res = &mut worker => {
tracing::error!(result = ?res, "analysis worker stopped unexpectedly; daemon exiting for restart");
break false;
}
}
};
shutdown_listeners(shutdown.energy, shutdown.listeners);
if !graceful {
return Err(super::DaemonError::AnalysisWorkerStopped);
}
drain_to_worker_and_join(window, energy_sources, work_tx, worker, metrics).await;
Ok(())
}
async fn run_analysis_worker(mut work_rx: mpsc::Receiver<AnalysisBatch>, wctx: AnalysisWorkerCtx) {
let mut db_waste_sticky: Option<(DatabaseWaste, u64)> = None;
let mut msg_waste_sticky: Option<(MessagingWaste, u64)> = None;
while let Some(batch) = work_rx.recv().await {
wctx.metrics.analysis_queue_depth.dec();
process_traces(
batch.traces,
ProcessTracesCtx {
detect_config: &wctx.detect_config,
green_enabled: wctx.green_enabled,
carbon_ctx: batch.carbon_ctx.as_ref(),
metrics: &wctx.metrics,
confidence: wctx.confidence,
findings_store: &wctx.findings_store,
traces_store: &wctx.traces_store,
correlator: wctx.correlator.as_deref(),
green_summary_cell: &wctx.green_summary_cell,
archive_tx: wctx.archive_tx.as_ref(),
db_waste_sticky: &mut db_waste_sticky,
msg_waste_sticky: &mut msg_waste_sticky,
waste_sticky_ttl_ms: wctx.waste_sticky_ttl_ms,
},
)
.await;
}
}
pub(crate) const MAX_SERVICE_CARDINALITY: usize = 1024;
struct ServiceMeter {
known_services: std::collections::HashMap<String, prometheus::Counter>,
max_service_cardinality: usize,
service_cap_warned: bool,
}
impl ServiceMeter {
fn record(&mut self, service: &str, metrics: &MetricsState) {
if let Some(child) = self.known_services.get(service) {
child.inc();
} else if self.known_services.len() < self.max_service_cardinality {
let child = metrics.service_io_ops_total.with_label_values(&[service]);
child.inc();
self.known_services.insert(service.to_string(), child);
} else {
metrics.service_io_ops_overflow_total.inc();
if !self.service_cap_warned {
tracing::warn!(
cap = self.max_service_cardinality,
"Service cardinality cap reached; new services will \
not have per-service I/O op counters"
);
self.service_cap_warned = true;
}
}
}
}
async fn ingest_event_batch(
events: Vec<SpanEvent>,
sampling_rate: f64,
window: &Arc<Mutex<TraceWindow>>,
metrics: &MetricsState,
service_meter: &mut ServiceMeter,
) -> Vec<(String, Vec<normalize::NormalizedEvent>)> {
let events = apply_sampling(events, sampling_rate);
let event_count = events.len();
let normalized: Vec<_> = events.into_iter().map(normalize::normalize).collect();
for event in &normalized {
service_meter.record(event.event.service.as_ref(), metrics);
}
let now_ms = current_time_ms();
let mut lru_evicted = Vec::new();
{
let mut w = window.lock().await;
for event in normalized {
if let Some(evicted) = w.push(event, now_ms) {
lru_evicted.push(evicted);
}
}
metrics.active_traces.set(w.active_traces() as f64);
}
metrics.events_processed_total.inc_by(event_count as f64);
lru_evicted
}
async fn evict_expired_traces(
window: &Arc<Mutex<TraceWindow>>,
metrics: &MetricsState,
) -> Vec<(String, Vec<normalize::NormalizedEvent>)> {
let now_ms = current_time_ms();
let mut w = window.lock().await;
let expired = w.evict_expired(now_ms);
metrics.active_traces.set(w.active_traces() as f64);
expired
}
fn build_owned_tick_ctx(sources: &EnergySources<'_>) -> Arc<score::carbon::CarbonContext> {
match build_tick_ctx(sources, score::scaphandre::monotonic_ms()) {
std::borrow::Cow::Borrowed(_) => Arc::clone(&sources.base_carbon_ctx),
std::borrow::Cow::Owned(ctx) => Arc::new(ctx),
}
}
fn enqueue_for_analysis(
traces: Vec<(String, Vec<normalize::NormalizedEvent>)>,
sources: &EnergySources<'_>,
work_tx: &mpsc::Sender<AnalysisBatch>,
metrics: &MetricsState,
) {
if traces.is_empty() {
return;
}
let trace_count = traces.len();
match work_tx.try_reserve() {
Ok(permit) => {
metrics.analysis_queue_depth.inc();
permit.send(AnalysisBatch::new(traces, sources));
}
Err(mpsc::error::TrySendError::Full(())) => {
metrics.record_shed(trace_count);
tracing::warn!(traces = trace_count, "analysis queue full, shedding batch");
}
Err(mpsc::error::TrySendError::Closed(())) => {
metrics.record_shed(trace_count);
tracing::error!(
traces = trace_count,
"analysis worker stopped, shedding batch"
);
}
}
}
async fn drain_to_worker_and_join(
window: &Arc<Mutex<TraceWindow>>,
sources: &EnergySources<'_>,
work_tx: mpsc::Sender<AnalysisBatch>,
worker: tokio::task::JoinHandle<()>,
metrics: &MetricsState,
) {
let remaining = {
let mut w = window.lock().await;
w.drain_all()
};
if !remaining.is_empty() {
let trace_count = remaining.len();
let batch = AnalysisBatch::new(remaining, sources);
if work_tx.send(batch).await.is_ok() {
metrics.analysis_queue_depth.inc();
} else {
metrics.record_shed(trace_count);
tracing::error!(
traces = trace_count,
"analysis worker stopped before shutdown drain"
);
}
}
drop(work_tx);
let _ = worker.await;
}
fn shutdown_listeners(energy: EnergyScraperHandles<'_>, listeners: ListenerHandles<'_>) {
if let Some(handle) = energy.emaps {
handle.abort();
}
if let Some(handle) = energy.cloud {
handle.abort();
}
if let Some(handle) = energy.redfish {
handle.abort();
}
if let Some(handle) = energy.kepler {
handle.abort();
}
if let Some(handle) = energy.scaphandre {
handle.abort();
}
if let Some(handle) = energy.alumet {
handle.abort();
}
listeners.grpc.abort();
listeners.http.abort();
if let Some(handle) = listeners.json_socket {
handle.abort();
}
}
fn build_tick_ctx<'s>(
sources: &'s EnergySources<'_>,
now: u64,
) -> std::borrow::Cow<'s, score::carbon::CarbonContext> {
let base = &*sources.base_carbon_ctx;
let EnergySources {
alumet_state,
alumet_db_state,
alumet_broker_state,
static_broker,
alumet_staleness_ms,
scaphandre_state,
scaphandre_staleness_ms,
kepler_state,
kepler_staleness_ms,
redfish_state,
redfish_staleness_ms,
cloud_state,
cloud_staleness_ms,
emaps_state,
emaps_staleness_ms,
..
} = *sources;
let cloud_snap = cloud_state
.map(|s| s.snapshot(now, cloud_staleness_ms))
.unwrap_or_default();
let redfish_snap = redfish_state
.map(|s| s.snapshot(now, redfish_staleness_ms))
.unwrap_or_default();
let kepler_snap = kepler_state
.map(|s| s.snapshot(now, kepler_staleness_ms))
.unwrap_or_default();
let scaph_snap = scaphandre_state
.map(|s| s.snapshot(now, scaphandre_staleness_ms))
.unwrap_or_default();
let alumet_snap = alumet_state
.map(|s| s.snapshot(now, alumet_staleness_ms))
.unwrap_or_default();
let emaps_snap = emaps_state
.map(|s| s.snapshot_with_metadata(now, emaps_staleness_ms))
.unwrap_or_default();
let db_window_kwh = alumet_db_state.and_then(|db| db.take_window_kwh(now, alumet_staleness_ms));
let (measured_broker_kwh, declared_broker_kwh) = take_broker_energy(
alumet_broker_state,
static_broker.map(|(_, state)| state),
now,
alumet_staleness_ms,
);
if cloud_snap.is_empty()
&& redfish_snap.is_empty()
&& kepler_snap.is_empty()
&& scaph_snap.is_empty()
&& alumet_snap.is_empty()
&& emaps_snap.is_empty()
&& db_window_kwh.is_none()
&& measured_broker_kwh.is_none()
&& declared_broker_kwh.is_none()
{
return std::borrow::Cow::Borrowed(base);
}
let mut merged: std::collections::HashMap<String, score::carbon::EnergyEntry> =
std::collections::HashMap::with_capacity(
cloud_snap.len()
+ redfish_snap.len()
+ kepler_snap.len()
+ scaph_snap.len()
+ alumet_snap.len(),
);
for (service, energy_kwh) in cloud_snap {
merged.insert(service, score::carbon::EnergyEntry::cloud(energy_kwh));
}
for (service, energy_kwh) in redfish_snap {
merged.insert(service, score::carbon::EnergyEntry::redfish(energy_kwh));
}
for (service, energy_kwh) in kepler_snap {
merged.insert(service, score::carbon::EnergyEntry::kepler(energy_kwh));
}
for (service, energy_kwh) in scaph_snap {
merged.insert(service, score::carbon::EnergyEntry::scaphandre(energy_kwh));
}
for (service, energy_kwh) in alumet_snap {
merged.insert(service, score::carbon::EnergyEntry::alumet(energy_kwh));
}
let mut ctx = base.clone();
ctx.energy_snapshot = if merged.is_empty() {
None
} else {
Some(merged)
};
if !emaps_snap.is_empty() {
ctx.real_time_intensity = Some(emaps_snap);
}
if let (Some(kwh), Some(db)) = (db_window_kwh, ctx.db_energy.as_mut()) {
db.window_kwh = kwh;
}
if let Some(broker) = ctx.broker_energy.as_mut() {
patch_broker_energy(
broker,
measured_broker_kwh,
declared_broker_kwh.zip(static_broker.map(|(cfg, _)| cfg)),
);
}
std::borrow::Cow::Owned(ctx)
}
fn take_broker_energy(
alumet_state: Option<&DbEnergyState>,
declared: Option<&score::broker_static::StaticBrokerState>,
now: u64,
alumet_staleness_ms: u64,
) -> (Option<f64>, Option<f64>) {
let measured_owns_the_timeline =
alumet_state.is_some_and(|b| b.has_recent_sample(now, alumet_staleness_ms));
if !measured_owns_the_timeline {
if declared.is_some_and(score::broker_static::StaticBrokerState::outage_billed) {
if let Some(state) = alumet_state {
state.discard_pending();
}
} else if let Some(kwh) =
alumet_state.and_then(|b| b.take_window_kwh(now, alumet_staleness_ms))
{
if let Some(state) = declared {
state.take_window_kwh(now);
}
return (Some(kwh), None);
}
let declared_kwh = declared.and_then(|state| state.take_window_kwh(now));
if declared_kwh.is_some()
&& let Some(state) = declared
{
state.mark_outage_billed();
}
return (None, declared_kwh);
}
if declared.is_some_and(score::broker_static::StaticBrokerState::clear_outage_billed)
&& let Some(state) = alumet_state
{
state.discard_pending();
}
let measured = alumet_state.and_then(|b| b.take_window_kwh(now, alumet_staleness_ms));
if let Some(state) = declared {
state.take_window_kwh(now);
}
(measured, None)
}
fn patch_broker_energy(
broker: &mut score::carbon::DbEnergyContext,
measured_kwh: Option<f64>,
declared: Option<(f64, &score::broker_static::StaticBrokerConfig)>,
) {
if let Some(kwh) = measured_kwh {
broker.window_kwh = kwh;
broker.model = score::carbon::CO2_MODEL_ALUMET;
} else if let Some((kwh, cfg)) = declared {
broker.window_kwh = kwh;
broker.model = crate::report::BROKER_WASTE_MODEL_SPECPOWER;
broker.region.clone_from(&cfg.region);
}
}
fn record_slow_durations(traces: &[Trace], detect_config: &DetectConfig, metrics: &MetricsState) {
let slow_threshold_us = detect_config.slow_threshold_ms.saturating_mul(1000);
let hist_sql = metrics.slow_duration_seconds.with_label_values(&["sql"]);
let hist_http = metrics
.slow_duration_seconds
.with_label_values(&["http_out"]);
let hist_messaging = metrics
.slow_duration_seconds
.with_label_values(&["messaging"]);
for trace in traces {
for span in &trace.spans {
if span.event.duration_us > slow_threshold_us {
let hist = match span.event.event_type {
crate::event::EventType::Sql => &hist_sql,
crate::event::EventType::HttpOut => &hist_http,
crate::event::EventType::Messaging => &hist_messaging,
};
hist.observe(span.event.duration_us as f64 / 1_000_000.0);
}
}
}
}
fn emit_findings_and_update_metrics(
trace_count: usize,
findings: &[detect::Finding],
green_summary: &GreenSummary,
metrics: &MetricsState,
) {
use std::io::Write;
metrics.traces_analyzed_total.inc_by(trace_count as f64);
metrics
.total_io_ops
.inc_by(green_summary.total_io_ops as f64);
metrics
.avoidable_io_ops
.inc_by(green_summary.avoidable_io_ops as f64);
let cumulative_total = metrics.total_io_ops.get();
if cumulative_total > 0.0 {
metrics
.io_waste_ratio
.set(metrics.avoidable_io_ops.get() / cumulative_total);
}
metrics.energy_kwh.set(green_summary.energy_kwh);
metrics
.carbon_gco2
.set(green_summary.regions.iter().map(|r| r.co2_gco2).sum());
metrics.record_exemplars(findings, green_summary);
let stdout = std::io::stdout();
let mut lock = stdout.lock();
for finding in findings {
metrics
.findings_total
.with_label_values(&[finding.finding_type.as_str(), finding.severity.as_str()])
.inc();
if serde_json::to_writer(&mut lock, finding).is_ok() {
let _ = writeln!(lock);
}
}
}
struct ProcessTracesCtx<'a> {
detect_config: &'a DetectConfig,
green_enabled: bool,
carbon_ctx: &'a score::carbon::CarbonContext,
metrics: &'a MetricsState,
confidence: Confidence,
findings_store: &'a findings_store::FindingsStore,
traces_store: &'a super::traces_store::TracesStore,
correlator: Option<&'a Mutex<detect::correlate_cross::CrossTraceCorrelator>>,
green_summary_cell: &'a Arc<RwLock<GreenSummary>>,
archive_tx: Option<&'a mpsc::Sender<super::archive::OwnedArchive>>,
db_waste_sticky: &'a mut Option<(DatabaseWaste, u64)>,
msg_waste_sticky: &'a mut Option<(MessagingWaste, u64)>,
waste_sticky_ttl_ms: u64,
}
async fn publish_live_summary(green_summary: &GreenSummary, ctx: &mut ProcessTracesCtx<'_>) {
let now_ms = current_time_ms();
let restored = sticky_waste_figure(
green_summary.database_waste.as_ref(),
ctx.db_waste_sticky,
now_ms,
ctx.waste_sticky_ttl_ms,
);
let restored_msg = sticky_waste_figure(
green_summary.messaging_waste.as_ref(),
ctx.msg_waste_sticky,
now_ms,
ctx.waste_sticky_ttl_ms,
);
let mut cell = ctx.green_summary_cell.write().await;
cell.clone_from(green_summary);
cell.database_waste = restored;
cell.messaging_waste = restored_msg;
}
fn sticky_waste_figure<T: Clone>(
fresh: Option<&T>,
sticky: &mut Option<(T, u64)>,
now_ms: u64,
ttl_ms: u64,
) -> Option<T> {
if let Some(figure) = fresh {
*sticky = Some((figure.clone(), now_ms));
return Some(figure.clone());
}
match sticky {
Some((figure, at)) if now_ms.saturating_sub(*at) <= ttl_ms && ttl_ms > 0 => {
Some(figure.clone())
}
_ => {
*sticky = None;
None
}
}
}
async fn process_traces(
traces: Vec<(String, Vec<normalize::NormalizedEvent>)>,
mut ctx: ProcessTracesCtx<'_>,
) {
if traces.is_empty() {
return;
}
let trace_count = traces.len();
let trace_structs: Vec<Trace> = traces
.into_iter()
.map(|(trace_id, spans)| Trace { trace_id, spans })
.collect();
let findings = detect::run_full_detection(&trace_structs, ctx.detect_config);
record_slow_durations(&trace_structs, ctx.detect_config, ctx.metrics);
let (mut findings, green_summary, per_endpoint_io_ops) = if ctx.green_enabled {
score::score_green(&trace_structs, findings, Some(ctx.carbon_ctx))
} else {
let total_io_ops = trace_structs.iter().map(|t| t.spans.len()).sum();
(findings, GreenSummary::disabled(total_io_ops), Vec::new())
};
publish_live_summary(&green_summary, &mut ctx).await;
detect::apply_confidence(&mut findings, ctx.confidence);
crate::acknowledgments::enrich_with_signatures(&mut findings);
let findings = findings;
let now_ms = current_time_ms();
if !findings.is_empty() {
ctx.findings_store.push_batch(&findings, now_ms).await;
ctx.traces_store.retain_for(&trace_structs, &findings).await;
#[allow(clippy::cast_precision_loss)] ctx.metrics
.stored_findings
.set(ctx.findings_store.len().await as f64);
}
if let Some(correlator) = ctx.correlator {
let evicted = correlator.lock().await.ingest(&findings, now_ms);
if evicted > 0 {
static CAP_WARNED: std::sync::atomic::AtomicBool =
std::sync::atomic::AtomicBool::new(false);
ctx.metrics
.correlator_pairs_evicted_total
.inc_by(evicted as u64);
if !CAP_WARNED.swap(true, std::sync::atomic::Ordering::Relaxed) {
tracing::warn!(
evicted,
"correlator pair cap reached, dropping pairs (see \
perf_sentinel_correlator_pairs_evicted_total)"
);
}
}
}
emit_findings_and_update_metrics(trace_count, &findings, &green_summary, ctx.metrics);
if let Some(archive_tx) = ctx.archive_tx {
let events_processed = trace_structs.iter().map(|t| t.spans.len()).sum();
let disclosure_waste = green_summary.co2.is_some().then(|| {
score::canonical::compute_disclosure_waste(
&trace_structs,
&green_summary,
ctx.detect_config,
)
});
let report = crate::report::Report {
analysis: crate::report::Analysis {
duration_ms: 0,
events_processed,
traces_analyzed: trace_count,
ingest: None,
},
findings,
green_summary,
quality_gate: crate::report::QualityGate {
passed: true,
rules: vec![],
},
per_endpoint_io_ops,
correlations: vec![],
embedded_traces: vec![],
warnings: vec![],
warning_details: vec![],
acknowledged_findings: vec![],
binary_version: env!("CARGO_PKG_VERSION").to_string(),
detection_config: Some(ctx.detect_config.clone()),
disclosure_waste,
};
let archive = super::archive::OwnedArchive {
ts: chrono::Utc::now(),
report,
};
super::archive::try_send(archive_tx, archive);
}
}
fn current_time_ms() -> u64 {
if let Ok(duration) = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH) {
u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
} else {
tracing::warn!(
"System clock is before Unix epoch; using 0 as current_time_ms. \
Check system time configuration."
);
0
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use super::*;
use crate::correlate::window::WindowConfig;
use crate::event::{EventSource, EventType, SpanEvent};
use core::assert_matches;
fn make_normalized(trace_id: &str, target: &str) -> normalize::NormalizedEvent {
make_normalized_for_service(trace_id, "test", target)
}
fn make_normalized_for_service(
trace_id: &str,
service: &str,
target: &str,
) -> normalize::NormalizedEvent {
let mut event = crate::test_helpers::make_sql_event_with_duration(
trace_id,
"s1",
target,
"2025-07-10T14:32:01.123Z",
100,
);
event.service = Arc::from(service);
normalize::normalize(event)
}
fn default_detect_config() -> DetectConfig {
DetectConfig {
n_plus_one_threshold: 5,
window_ms: 500,
slow_threshold_ms: 500,
slow_min_occurrences: 3,
max_fanout: 20,
chatty_service_min_calls: 15,
pool_saturation_concurrent_threshold: 10,
serialized_min_sequential: 3,
sanitizer_aware_classification: SanitizerAwareMode::default(),
}
}
fn empty_carbon_ctx() -> score::carbon::CarbonContext {
score::carbon::CarbonContext::default()
}
fn noop_traces_store() -> &'static crate::daemon::traces_store::TracesStore {
static STORE: std::sync::OnceLock<crate::daemon::traces_store::TracesStore> =
std::sync::OnceLock::new();
STORE.get_or_init(|| crate::daemon::traces_store::TracesStore::new(0, 0))
}
fn test_ctx<'a>(
detect_config: &'a DetectConfig,
carbon_ctx: &'a score::carbon::CarbonContext,
metrics: &'a MetricsState,
findings_store: &'a findings_store::FindingsStore,
green_enabled: bool,
green_summary_cell: &'a Arc<RwLock<GreenSummary>>,
) -> ProcessTracesCtx<'a> {
ProcessTracesCtx {
detect_config,
traces_store: noop_traces_store(),
green_enabled,
carbon_ctx,
metrics,
confidence: Confidence::DaemonStaging,
findings_store,
correlator: None,
green_summary_cell,
archive_tx: None,
db_waste_sticky: Box::leak(Box::new(None)),
msg_waste_sticky: Box::leak(Box::new(None)),
waste_sticky_ttl_ms: 0,
}
}
fn fresh_green_cell() -> Arc<RwLock<GreenSummary>> {
Arc::new(RwLock::new(GreenSummary::disabled(0)))
}
#[tokio::test]
async fn process_traces_empty_does_nothing() {
let metrics = MetricsState::new();
let ctx = empty_carbon_ctx();
let store = findings_store::FindingsStore::new(100);
let detect_config = default_detect_config();
let cell = fresh_green_cell();
process_traces(
vec![],
test_ctx(&detect_config, &ctx, &metrics, &store, true, &cell),
)
.await;
}
#[tokio::test]
async fn process_traces_with_n_plus_one() {
let events: Vec<_> = (1..=6)
.map(|i| {
make_normalized(
"t1",
&format!("SELECT * FROM order_item WHERE order_id = {i}"),
)
})
.collect();
let metrics = MetricsState::new();
let ctx = empty_carbon_ctx();
let store = findings_store::FindingsStore::new(100);
let detect_config = default_detect_config();
let cell = fresh_green_cell();
process_traces(
vec![("t1".to_string(), events)],
test_ctx(&detect_config, &ctx, &metrics, &store, true, &cell),
)
.await;
}
#[tokio::test]
async fn process_traces_clean_no_finding() {
let events = vec![
make_normalized("t1", "SELECT * FROM users WHERE id = 1"),
make_normalized("t1", "SELECT * FROM orders WHERE id = 2"),
];
let metrics = MetricsState::new();
let ctx = empty_carbon_ctx();
let store = findings_store::FindingsStore::new(100);
let detect_config = default_detect_config();
let cell = fresh_green_cell();
process_traces(
vec![("t1".to_string(), events)],
test_ctx(&detect_config, &ctx, &metrics, &store, true, &cell),
)
.await;
}
#[test]
fn current_time_ms_returns_nonzero() {
let ms = current_time_ms();
assert!(ms > 0, "current_time_ms should return a positive value");
}
#[test]
fn evict_expired_returns_traces() {
let config = WindowConfig {
trace_ttl_ms: 100,
..Default::default()
};
let mut w = TraceWindow::new(config);
let event = normalize::normalize(SpanEvent {
timestamp: "2025-07-10T14:32:01.123Z".to_string(),
trace_id: "t1".to_string(),
span_id: "s1".to_string(),
parent_span_id: None,
link_trace_id: None,
service: Arc::from("test"),
grouping: Vec::new(),
cloud_region: None,
event_type: EventType::Sql,
operation: "SELECT".to_string(),
target: "SELECT 1".to_string(),
duration_us: 100,
source: EventSource {
endpoint: "GET /test".to_string(),
method: "Test::test".to_string(),
},
status_code: None,
response_size_bytes: None,
code_function: None,
code_filepath: None,
code_lineno: None,
code_namespace: None,
instrumentation_scopes: Vec::new(),
});
w.push(event, 0);
assert_eq!(w.active_traces(), 1);
let expired = w.evict_expired(50);
assert!(expired.is_empty());
assert_eq!(w.active_traces(), 1);
let expired = w.evict_expired(150);
assert_eq!(expired.len(), 1);
assert_eq!(expired[0].0, "t1");
assert_eq!(expired[0].1.len(), 1);
assert_eq!(w.active_traces(), 0);
}
#[tokio::test]
async fn process_traces_updates_metrics() {
let events: Vec<_> = (1..=6)
.map(|i| {
make_normalized(
"t1",
&format!("SELECT * FROM order_item WHERE order_id = {i}"),
)
})
.collect();
let metrics = MetricsState::new();
let ctx = empty_carbon_ctx();
let store = findings_store::FindingsStore::new(100);
let detect_config = default_detect_config();
let cell = fresh_green_cell();
process_traces(
vec![("t1".to_string(), events)],
test_ctx(&detect_config, &ctx, &metrics, &store, true, &cell),
)
.await;
let output = metrics.render();
assert!(output.contains("perf_sentinel_traces_analyzed_total"));
assert!(output.contains("perf_sentinel_findings_total"));
}
#[tokio::test]
async fn process_traces_green_disabled() {
let events: Vec<_> = (1..=6)
.map(|i| {
make_normalized(
"t1",
&format!("SELECT * FROM order_item WHERE order_id = {i}"),
)
})
.collect();
let metrics = MetricsState::new();
let ctx = empty_carbon_ctx();
let store = findings_store::FindingsStore::new(100);
let detect_config = default_detect_config();
let cell = fresh_green_cell();
process_traces(
vec![("t1".to_string(), events)],
test_ctx(&detect_config, &ctx, &metrics, &store, false, &cell),
)
.await;
assert!((metrics.avoidable_io_ops.get() - 0.0).abs() < f64::EPSILON);
assert!(metrics.total_io_ops.get() > 0.0);
}
#[tokio::test]
async fn process_traces_publishes_green_summary_to_cell() {
let events: Vec<_> = (1..=6)
.map(|i| {
make_normalized(
"t1",
&format!("SELECT * FROM order_item WHERE order_id = {i}"),
)
})
.collect();
let metrics = MetricsState::new();
let ctx = empty_carbon_ctx();
let store = findings_store::FindingsStore::new(100);
let detect_config = default_detect_config();
let cell = fresh_green_cell();
process_traces(
vec![("t1".to_string(), events)],
test_ctx(&detect_config, &ctx, &metrics, &store, true, &cell),
)
.await;
let snapshot = cell.read().await.clone();
assert!(snapshot.total_io_ops > 0, "cell should reflect the batch");
}
#[test]
fn build_tick_ctx_no_scrapers_yields_borrowed_cow() {
let base = Arc::new(score::carbon::CarbonContext::default());
let sources = no_scrapers(&base);
let ctx = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
assert_matches!(ctx, std::borrow::Cow::Borrowed(_));
assert!(ctx.energy_snapshot.is_none());
}
#[test]
fn build_tick_ctx_scaphandre_only() {
let base = Arc::new(score::carbon::CarbonContext::default());
let scaph = ScaphandreState::new();
scaph.insert_for_test("svc-a".into(), 1e-7, 100);
let mut sources = no_scrapers(&base);
sources.scaphandre_state = Some(&scaph);
sources.scaphandre_staleness_ms = 500;
let ctx = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
let snap = ctx.energy_snapshot.as_ref().unwrap();
assert_eq!(snap.len(), 1);
assert_eq!(snap["svc-a"].model_tag, "scaphandre_rapl");
}
#[test]
fn build_tick_ctx_cloud_only() {
let base = Arc::new(score::carbon::CarbonContext::default());
let cloud = CloudEnergyState::new();
cloud.insert_for_test("svc-b".into(), 2e-7, 100);
let mut sources = no_scrapers(&base);
sources.cloud_state = Some(&cloud);
sources.cloud_staleness_ms = 500;
let ctx = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
let snap = ctx.energy_snapshot.as_ref().unwrap();
assert_eq!(snap.len(), 1);
assert_eq!(snap["svc-b"].model_tag, "cloud_specpower");
}
#[test]
fn build_tick_ctx_kepler_only() {
let base = Arc::new(score::carbon::CarbonContext::default());
let kepler = KeplerState::new();
kepler.insert_for_test("svc-k".into(), 4e-7, 100);
let mut sources = no_scrapers(&base);
sources.kepler_state = Some(&kepler);
sources.kepler_staleness_ms = 500;
let ctx = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
let snap = ctx.energy_snapshot.as_ref().unwrap();
assert_eq!(snap.len(), 1);
assert_eq!(snap["svc-k"].model_tag, "kepler_ebpf");
}
#[test]
fn build_tick_ctx_redfish_only() {
let base = Arc::new(score::carbon::CarbonContext::default());
let redfish = RedfishState::new();
redfish.insert_for_test("svc-r".into(), 6e-7, 100);
let mut sources = no_scrapers(&base);
sources.redfish_state = Some(&redfish);
sources.redfish_staleness_ms = 500;
let ctx = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
let snap = ctx.energy_snapshot.as_ref().unwrap();
assert_eq!(snap.len(), 1);
assert_eq!(snap["svc-r"].model_tag, "redfish_bmc");
}
#[test]
fn build_tick_ctx_scaphandre_overrides_kepler_overrides_cloud_for_same_service() {
let base = Arc::new(score::carbon::CarbonContext::default());
let scaph = ScaphandreState::new();
scaph.insert_for_test("svc-a".into(), 1e-7, 100);
let kepler = KeplerState::new();
kepler.insert_for_test("svc-a".into(), 2e-7, 100);
kepler.insert_for_test("svc-k".into(), 4e-7, 100);
let cloud = CloudEnergyState::new();
cloud.insert_for_test("svc-a".into(), 5e-7, 100);
cloud.insert_for_test("svc-b".into(), 3e-7, 100);
let mut sources = no_scrapers(&base);
sources.scaphandre_state = Some(&scaph);
sources.scaphandre_staleness_ms = 500;
sources.kepler_state = Some(&kepler);
sources.kepler_staleness_ms = 500;
sources.cloud_state = Some(&cloud);
sources.cloud_staleness_ms = 500;
let ctx = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
let snap = ctx.energy_snapshot.as_ref().unwrap();
assert_eq!(snap.len(), 3);
assert_eq!(snap["svc-a"].model_tag, "scaphandre_rapl");
assert!((snap["svc-a"].energy_per_op_kwh - 1e-7).abs() < 1e-15);
assert_eq!(snap["svc-k"].model_tag, "kepler_ebpf");
assert_eq!(snap["svc-b"].model_tag, "cloud_specpower");
}
#[test]
fn build_tick_ctx_alumet_only() {
let base = Arc::new(score::carbon::CarbonContext::default());
let alumet = AlumetState::new();
alumet.insert_for_test("svc-al".into(), 8e-7, 100);
let mut sources = no_scrapers(&base);
sources.alumet_state = Some(&alumet);
sources.alumet_staleness_ms = 500;
let ctx = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
let snap = ctx.energy_snapshot.as_ref().unwrap();
assert_eq!(snap.len(), 1);
assert_eq!(snap["svc-al"].model_tag, "alumet_rapl");
}
fn waste_fixture(ratio: f64) -> DatabaseWaste {
DatabaseWaste {
energy_kwh: 0.01,
waste_kwh: 0.01 * ratio,
waste_gco2: None,
energy_gco2: None,
region: None,
sql_waste_ratio: ratio,
model: "alumet_rapl".to_string(),
}
}
#[test]
fn sticky_waste_figure_bridges_gaps_then_ages_out() {
let mut sticky = None;
let fresh = waste_fixture(0.4);
let out = sticky_waste_figure(Some(&fresh), &mut sticky, 1_000, 30_000);
assert_eq!(out.as_ref(), Some(&fresh));
let out = sticky_waste_figure(None, &mut sticky, 10_000, 30_000);
assert_eq!(out.as_ref(), Some(&fresh));
let out = sticky_waste_figure(None, &mut sticky, 40_000, 30_000);
assert!(out.is_none());
assert!(sticky.is_none(), "aged-out figure must be dropped");
}
#[test]
fn sticky_waste_figure_disabled_at_zero_ttl() {
let mut sticky = None;
let fresh = waste_fixture(0.2);
assert!(sticky_waste_figure(Some(&fresh), &mut sticky, 1_000, 0).is_some());
assert!(sticky_waste_figure(None, &mut sticky, 1_001, 0).is_none());
}
#[test]
fn build_tick_ctx_database_energy_forces_owned_then_consumes() {
let base = Arc::new(score::carbon::CarbonContext {
db_energy: Some(score::carbon::DbEnergyContext {
window_kwh: 0.0,
region: None,
..Default::default()
}),
..score::carbon::CarbonContext::default()
});
let db = DbEnergyState::new();
let now = score::scaphandre::monotonic_ms();
db.add_window_kwh(2e-6, now);
let mut sources = no_scrapers(&base);
sources.alumet_db_state = Some(&db);
sources.alumet_staleness_ms = 60_000;
let ctx = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
assert!(
matches!(ctx, std::borrow::Cow::Owned(_)),
"fresh db energy must not take the borrowed fast path"
);
let kwh = ctx.db_energy.as_ref().unwrap().window_kwh;
assert!((kwh - 2e-6).abs() < 1e-18);
let ctx2 = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
assert!(matches!(ctx2, std::borrow::Cow::Borrowed(_)));
assert!((ctx2.db_energy.as_ref().unwrap().window_kwh - 0.0).abs() < f64::EPSILON);
}
#[test]
fn build_tick_ctx_broker_only_energy_forces_owned() {
let base = Arc::new(score::carbon::CarbonContext {
broker_energy: Some(score::carbon::DbEnergyContext {
window_kwh: 0.0,
region: None,
..Default::default()
}),
..score::carbon::CarbonContext::default()
});
let broker = DbEnergyState::new();
broker.add_window_kwh(3e-6, 10_000);
let mut sources = no_scrapers(&base);
sources.alumet_broker_state = Some(&broker);
sources.alumet_staleness_ms = 60_000;
let ctx = build_tick_ctx(&sources, 10_000);
assert!(
matches!(ctx, std::borrow::Cow::Owned(_)),
"fresh broker energy must not take the borrowed fast path"
);
let kwh = ctx.broker_energy.as_ref().unwrap().window_kwh;
assert!((kwh - 3e-6).abs() < 1e-18);
let ctx2 = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
assert!(matches!(ctx2, std::borrow::Cow::Borrowed(_)));
}
fn declared_cfg(nodes: u32) -> score::broker_static::StaticBrokerConfig {
score::broker_static::StaticBrokerConfig {
nodes,
instance_type: "m5.2xlarge".to_string(),
provider: "aws".to_string(),
region: Some("eu-west-3".to_string()),
}
}
#[test]
fn a_measured_broker_outranks_the_declared_cluster() {
let measured = DbEnergyState::new();
measured.add_window_kwh(5e-6, 10_000);
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
let (m, d) = take_broker_energy(Some(&measured), Some(&state), 10_000, 60_000);
assert_eq!(m, Some(5e-6));
assert!(
d.is_none(),
"a declaration must not be billed beside a measurement"
);
}
#[test]
fn a_gap_between_alumet_deltas_is_not_billed_by_the_declaration() {
let measured = DbEnergyState::new();
measured.add_window_kwh(5e-6, 10_000);
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
take_broker_energy(Some(&measured), Some(&state), 10_000, 60_000);
let (m, d) = take_broker_energy(Some(&measured), Some(&state), 20_000, 60_000);
assert!(m.is_none(), "no delta accumulated");
assert!(
d.is_none(),
"the declaration would re-bill what the next Alumet delta covers"
);
}
#[test]
fn a_stale_alumet_hands_the_window_over_to_the_declaration() {
let measured = DbEnergyState::new();
measured.add_window_kwh(5e-6, 1_000);
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
let (m, d) = take_broker_energy(Some(&measured), Some(&state), 100_000, 10_000);
assert!(m.is_none());
assert!(d.is_some_and(|k| k > 0.0));
}
#[test]
fn recovery_after_a_fallback_stretch_drops_the_banked_energy() {
let measured = DbEnergyState::new();
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
measured.add_window_kwh(5e-6, 1_000);
let (_, d) = take_broker_energy(Some(&measured), Some(&state), 100_000, 10_000);
assert!(d.is_some(), "the declaration covers the outage");
measured.add_window_kwh(2e-6, 101_000);
let (m, d2) = take_broker_energy(Some(&measured), Some(&state), 101_000, 10_000);
assert!(m.is_none(), "the recovery delta covers billed wall clock");
assert!(d2.is_none(), "the measurement owns the timeline again");
measured.add_window_kwh(3e-6, 102_000);
let (m2, _) = take_broker_energy(Some(&measured), Some(&state), 102_000, 10_000);
let delivered = m2.expect("the measurement resumes");
assert!(
(delivered - 3e-6).abs() < 1e-18,
"only the joules after the handover may be billed, got {delivered}"
);
}
#[test]
fn a_bank_landing_after_a_billed_outage_is_dropped_not_delivered() {
let measured = DbEnergyState::new();
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
measured.add_window_kwh(1e-6, 1_000);
take_broker_energy(Some(&measured), Some(&state), 1_000, 10_000);
measured.mark_alive(100_000);
let (_, d) = take_broker_energy(Some(&measured), Some(&state), 100_000, 10_000);
assert!(d.is_some(), "the declaration bills the outage");
measured.add_window_kwh(5e-6, 105_000);
measured.mark_alive(200_000);
let (m, _) = take_broker_energy(Some(&measured), Some(&state), 200_000, 10_000);
assert!(
m.is_none(),
"the banked delta covers wall clock the declaration already billed"
);
}
#[test]
fn sub_second_stale_ticks_do_not_erase_the_outage_marker() {
let measured = DbEnergyState::new();
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
measured.add_window_kwh(1e-6, 1_000);
take_broker_energy(Some(&measured), Some(&state), 1_000, 10_000);
measured.mark_alive(100_000);
let (_, d) = take_broker_energy(Some(&measured), Some(&state), 100_000, 10_000);
assert!(d.is_some(), "the declaration bills the outage");
for t in [100_300_u64, 100_600, 100_900] {
measured.mark_alive(t);
let (_, billed) = take_broker_energy(Some(&measured), Some(&state), t, 10_000);
assert!(billed.is_none(), "a sub-second tick bills nothing at t={t}");
}
measured.add_window_kwh(5e-6, 101_000);
measured.mark_alive(200_000);
let (m, _) = take_broker_energy(Some(&measured), Some(&state), 200_000, 10_000);
assert!(
m.is_none(),
"the marker must survive ticks that bill nothing"
);
}
#[test]
fn the_declared_marker_advances_while_alumet_owns_the_timeline() {
let measured = DbEnergyState::new();
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
for t in [10_000_u64, 20_000, 30_000] {
measured.add_window_kwh(1e-6, t);
take_broker_energy(Some(&measured), Some(&state), t, 60_000);
}
let (_, d) = take_broker_energy(Some(&measured), Some(&state), 100_000, 10_000);
let billed = d.expect("the fallback covers the outage");
let outage_kwh = cfg.cluster_watts() * 70_000.0 / 3_600_000.0 / 1000.0;
assert!(
(billed - outage_kwh).abs() < 1e-12,
"billed {billed} kWh, expected the 70 s outage alone"
);
}
#[test]
fn a_fallback_window_carries_the_declared_tag_and_region() {
let mut broker = score::carbon::DbEnergyContext {
window_kwh: 0.0,
region: Some("eu-west-1".to_string()),
model: score::carbon::CO2_MODEL_ALUMET,
};
let cfg = declared_cfg(1);
patch_broker_energy(&mut broker, None, Some((4.2e-6, &cfg)));
assert_eq!(
broker.model,
crate::report::BROKER_WASTE_MODEL_SPECPOWER,
"a fallback window must not be published as a measurement"
);
assert_eq!(broker.region.as_deref(), Some("eu-west-3"));
assert!((broker.window_kwh - 4.2e-6).abs() < 1e-18);
}
#[test]
fn a_measured_window_keeps_its_tag_when_both_sources_deliver() {
let mut broker = score::carbon::DbEnergyContext {
window_kwh: 0.0,
region: Some("eu-west-1".to_string()),
model: crate::report::BROKER_WASTE_MODEL_SPECPOWER,
};
let cfg = declared_cfg(1);
patch_broker_energy(&mut broker, Some(5e-6), Some((9e-6, &cfg)));
assert_eq!(broker.model, score::carbon::CO2_MODEL_ALUMET);
assert_eq!(broker.region.as_deref(), Some("eu-west-1"));
assert!((broker.window_kwh - 5e-6).abs() < 1e-18);
}
#[test]
fn build_tick_ctx_falls_back_to_the_declared_cluster() {
let base = Arc::new(score::carbon::CarbonContext {
broker_energy: Some(score::carbon::DbEnergyContext {
window_kwh: 0.0,
region: None,
..Default::default()
}),
..score::carbon::CarbonContext::default()
});
let declared = declared_cfg(3);
let declared_state = score::broker_static::StaticBrokerState::new(0, &declared);
let mut sources = no_scrapers(&base);
sources.static_broker = Some((&declared, &declared_state));
let ctx = build_tick_ctx(&sources, 60_000);
assert!(
matches!(ctx, std::borrow::Cow::Owned(_)),
"a declared cluster alone must leave the borrowed fast path"
);
let broker = ctx.broker_energy.as_ref().expect("broker context");
assert!(broker.window_kwh > 0.0);
assert_eq!(broker.model, crate::report::BROKER_WASTE_MODEL_SPECPOWER);
}
#[test]
fn build_tick_ctx_keeps_the_fast_path_on_a_sub_second_tick() {
let base = Arc::new(score::carbon::CarbonContext {
broker_energy: Some(score::carbon::DbEnergyContext::default()),
..score::carbon::CarbonContext::default()
});
let declared = declared_cfg(3);
let declared_state = score::broker_static::StaticBrokerState::new(0, &declared);
let mut sources = no_scrapers(&base);
sources.static_broker = Some((&declared, &declared_state));
let ctx = build_tick_ctx(&sources, 200);
assert!(matches!(ctx, std::borrow::Cow::Borrowed(_)));
}
#[test]
fn a_scrape_without_the_broker_label_hands_over_to_the_declaration() {
let measured = DbEnergyState::new();
measured.mark_alive(10_000);
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
let (m, d) = take_broker_energy(Some(&measured), Some(&state), 10_000, 60_000);
assert!(m.is_none(), "no sample ever carried the label");
assert!(
d.is_some_and(|k| k > 0.0),
"the declaration must cover a workload nothing measured"
);
}
#[test]
fn a_vanished_label_still_delivers_what_it_measured() {
let measured = DbEnergyState::new();
measured.add_window_kwh(4e-6, 10_000);
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
measured.mark_alive(100_000);
let (m, d) = take_broker_energy(Some(&measured), Some(&state), 100_000, 10_000);
assert_eq!(m, Some(4e-6), "banked measured energy must be delivered");
assert!(d.is_none(), "the declaration does not bill the same window");
let (m2, d2) = take_broker_energy(Some(&measured), Some(&state), 110_000, 10_000);
assert!(m2.is_none());
assert!(d2.is_some_and(|k| k > 0.0));
}
#[test]
fn an_unscraped_state_does_not_own_the_timeline_at_boot() {
let measured = DbEnergyState::new();
let cfg = declared_cfg(3);
let state = score::broker_static::StaticBrokerState::new(0, &cfg);
let (m, d) = take_broker_energy(Some(&measured), Some(&state), 5_000, 60_000);
assert!(m.is_none());
assert!(
d.is_some_and(|k| k > 0.0),
"a state that never saw a scrape must not suppress the fallback"
);
}
#[test]
fn build_tick_ctx_alumet_overrides_scaphandre_for_same_service() {
let base = Arc::new(score::carbon::CarbonContext::default());
let alumet = AlumetState::new();
alumet.insert_for_test("svc-a".into(), 1e-7, 100);
let scaph = ScaphandreState::new();
scaph.insert_for_test("svc-a".into(), 9e-7, 100);
scaph.insert_for_test("svc-s".into(), 3e-7, 100);
let mut sources = no_scrapers(&base);
sources.alumet_state = Some(&alumet);
sources.alumet_staleness_ms = 500;
sources.scaphandre_state = Some(&scaph);
sources.scaphandre_staleness_ms = 500;
let ctx = build_tick_ctx(&sources, score::scaphandre::monotonic_ms());
let snap = ctx.energy_snapshot.as_ref().unwrap();
assert_eq!(snap.len(), 2);
assert_eq!(snap["svc-a"].model_tag, "alumet_rapl");
assert!((snap["svc-a"].energy_per_op_kwh - 1e-7).abs() < 1e-15);
assert_eq!(snap["svc-s"].model_tag, "scaphandre_rapl");
}
#[test]
fn build_tick_ctx_stale_entries_filtered() {
let scaph = ScaphandreState::new();
scaph.insert_for_test("stale-svc".into(), 1e-7, 0);
let snap = scaph.snapshot(100, 1);
assert!(
snap.is_empty(),
"entry at time 0 should be stale when now=100, staleness=1"
);
scaph.insert_for_test("fresh-svc".into(), 2e-7, 99);
let snap2 = scaph.snapshot(100, 50);
assert!(snap2.contains_key("fresh-svc"));
assert!(!snap2.contains_key("stale-svc"));
}
fn no_scrapers(base: &Arc<score::carbon::CarbonContext>) -> EnergySources<'_> {
EnergySources {
base_carbon_ctx: base.clone(),
alumet_state: None,
alumet_db_state: None,
alumet_broker_state: None,
static_broker: None,
alumet_staleness_ms: 0,
scaphandre_state: None,
scaphandre_staleness_ms: 0,
kepler_state: None,
kepler_staleness_ms: 0,
redfish_state: None,
redfish_staleness_ms: 0,
cloud_state: None,
cloud_staleness_ms: 0,
emaps_state: None,
emaps_staleness_ms: 0,
}
}
fn one_trace_batch(id: &str) -> Vec<(String, Vec<normalize::NormalizedEvent>)> {
vec![(id.to_string(), vec![make_normalized(id, "SELECT 1")])]
}
fn test_window() -> Arc<Mutex<TraceWindow>> {
Arc::new(Mutex::new(TraceWindow::new(WindowConfig {
max_events_per_trace: 1000,
trace_ttl_ms: 30_000,
max_active_traces: std::num::NonZeroUsize::new(10_000).expect("nonzero"),
})))
}
fn test_worker_ctx(
metrics: &Arc<MetricsState>,
findings_store: &Arc<findings_store::FindingsStore>,
green_summary_cell: &Arc<RwLock<GreenSummary>>,
) -> AnalysisWorkerCtx {
AnalysisWorkerCtx {
detect_config: default_detect_config(),
traces_store: Arc::new(crate::daemon::traces_store::TracesStore::new(0, 0)),
green_enabled: true,
confidence: Confidence::DaemonStaging,
metrics: metrics.clone(),
findings_store: findings_store.clone(),
correlator: None,
green_summary_cell: green_summary_cell.clone(),
archive_tx: None,
waste_sticky_ttl_ms: 0,
}
}
#[tokio::test]
async fn ingestion_not_head_of_line_blocked_by_slow_analysis() {
let metrics = MetricsState::new();
let base = Arc::new(empty_carbon_ctx());
let sources = no_scrapers(&base);
let (work_tx, _work_rx) = mpsc::channel::<AnalysisBatch>(2);
for i in 0..10u32 {
enqueue_for_analysis(
one_trace_batch(&format!("t{i}")),
&sources,
&work_tx,
&metrics,
);
}
assert_eq!(metrics.analysis_queue_depth.get(), 2);
assert_eq!(metrics.analysis_shed_batches_total.get(), 8);
assert_eq!(metrics.analysis_shed_traces_total.get(), 8);
}
#[tokio::test]
async fn saturated_queue_sheds_and_increments_metric() {
let metrics = MetricsState::new();
let base = Arc::new(empty_carbon_ctx());
let sources = no_scrapers(&base);
let (work_tx, _work_rx) = mpsc::channel::<AnalysisBatch>(1);
enqueue_for_analysis(one_trace_batch("t1"), &sources, &work_tx, &metrics);
assert_eq!(metrics.analysis_queue_depth.get(), 1);
assert_eq!(metrics.analysis_shed_batches_total.get(), 0);
let batch = vec![
("t2".to_string(), vec![make_normalized("t2", "SELECT 1")]),
("t3".to_string(), vec![make_normalized("t3", "SELECT 1")]),
("t4".to_string(), vec![make_normalized("t4", "SELECT 1")]),
];
enqueue_for_analysis(batch, &sources, &work_tx, &metrics);
assert_eq!(metrics.analysis_shed_batches_total.get(), 1);
assert_eq!(metrics.analysis_shed_traces_total.get(), 3);
assert_eq!(metrics.analysis_queue_depth.get(), 1);
}
#[tokio::test]
async fn stopped_worker_counts_as_shed() {
let metrics = MetricsState::new();
let base = Arc::new(empty_carbon_ctx());
let sources = no_scrapers(&base);
let (work_tx, work_rx) = mpsc::channel::<AnalysisBatch>(4);
drop(work_rx);
let batch = vec![
("t1".to_string(), vec![make_normalized("t1", "SELECT 1")]),
("t2".to_string(), vec![make_normalized("t2", "SELECT 1")]),
];
enqueue_for_analysis(batch, &sources, &work_tx, &metrics);
assert_eq!(metrics.analysis_shed_batches_total.get(), 1);
assert_eq!(metrics.analysis_shed_traces_total.get(), 2);
assert_eq!(metrics.analysis_queue_depth.get(), 0);
}
#[tokio::test]
async fn correlator_pair_evictions_recorded_in_metrics() {
let metrics = MetricsState::new();
let carbon = empty_carbon_ctx();
let store = findings_store::FindingsStore::new(100);
let detect_config = default_detect_config();
let cell = fresh_green_cell();
let correlator = Mutex::new(detect::correlate_cross::CrossTraceCorrelator::new(
detect::correlate_cross::CorrelationConfig {
enabled: true,
max_tracked_pairs: 1,
lag_threshold_ms: 100_000,
min_co_occurrences: 1,
min_confidence: 0.0,
..Default::default()
},
));
let mut ctx = test_ctx(&detect_config, &carbon, &metrics, &store, true, &cell);
ctx.correlator = Some(&correlator);
let traces: Vec<_> = ["svc-a", "svc-b", "svc-c"]
.iter()
.enumerate()
.map(|(i, svc)| {
let trace_id = format!("t{i}");
let events: Vec<_> = (1..=6)
.map(|p| {
make_normalized_for_service(
&trace_id,
svc,
&format!("SELECT * FROM order_item WHERE order_id = {p}"),
)
})
.collect();
(trace_id, events)
})
.collect();
process_traces(traces, ctx).await;
assert!(
metrics.correlator_pairs_evicted_total.get() > 0,
"pair cap evictions must reach the metric"
);
}
#[test]
fn service_meter_overflow_counts_unattributed_ops() {
let metrics = MetricsState::new();
let mut meter = ServiceMeter {
known_services: std::collections::HashMap::new(),
max_service_cardinality: 2,
service_cap_warned: false,
};
for service in ["svc-a", "svc-b", "svc-c"] {
meter.record(service, &metrics);
meter.record(service, &metrics);
}
assert_eq!(metrics.service_io_ops_overflow_total.get(), 2);
for service in ["svc-a", "svc-b"] {
let count = metrics
.service_io_ops_total
.with_label_values(&[service])
.get();
assert!((count - 2.0).abs() < f64::EPSILON);
}
assert!(meter.service_cap_warned);
}
#[tokio::test]
async fn shed_traces_are_excluded_from_analysis_outputs() {
let metrics = Arc::new(MetricsState::new());
let store = Arc::new(findings_store::FindingsStore::new(100));
let cell = fresh_green_cell();
let base = Arc::new(empty_carbon_ctx());
let sources = no_scrapers(&base);
let (work_tx, work_rx) = mpsc::channel::<AnalysisBatch>(1);
let n_plus_one_events = |trace_id: &str| -> Vec<normalize::NormalizedEvent> {
(1..=6)
.map(|i| {
make_normalized(
trace_id,
&format!("SELECT * FROM order_item WHERE order_id = {i}"),
)
})
.collect()
};
enqueue_for_analysis(
vec![("kept".to_string(), n_plus_one_events("kept"))],
&sources,
&work_tx,
&metrics,
);
enqueue_for_analysis(
vec![("shed".to_string(), n_plus_one_events("shed"))],
&sources,
&work_tx,
&metrics,
);
assert_eq!(metrics.analysis_shed_batches_total.get(), 1);
assert_eq!(metrics.analysis_shed_traces_total.get(), 1);
let worker = tokio::spawn(run_analysis_worker(
work_rx,
test_worker_ctx(&metrics, &store, &cell),
));
drop(work_tx);
worker.await.expect("worker should drain and exit");
assert!((metrics.traces_analyzed_total.get() - 1.0).abs() < f64::EPSILON);
assert!(
!store.by_trace_id("kept").await.is_empty(),
"kept trace must reach the findings store"
);
assert!(
store.by_trace_id("shed").await.is_empty(),
"shed trace must never reach analysis outputs"
);
}
#[tokio::test]
async fn shutdown_drains_window_and_inflight_queue() {
let metrics = Arc::new(MetricsState::new());
let store = Arc::new(findings_store::FindingsStore::new(100));
let cell = fresh_green_cell();
let base = Arc::new(empty_carbon_ctx());
let sources = no_scrapers(&base);
let (work_tx, work_rx) = mpsc::channel::<AnalysisBatch>(4);
let worker = tokio::spawn(run_analysis_worker(
work_rx,
test_worker_ctx(&metrics, &store, &cell),
));
let inflight = vec![
("q1".to_string(), vec![make_normalized("q1", "SELECT 1")]),
("q2".to_string(), vec![make_normalized("q2", "SELECT 1")]),
];
enqueue_for_analysis(inflight, &sources, &work_tx, &metrics);
let window = test_window();
{
let mut w = window.lock().await;
for id in ["w1", "w2", "w3"] {
w.push(make_normalized(id, "SELECT 1"), 0);
}
}
drain_to_worker_and_join(&window, &sources, work_tx, worker, &metrics).await;
assert!((metrics.traces_analyzed_total.get() - 5.0).abs() < f64::EPSILON);
assert_eq!(metrics.analysis_queue_depth.get(), 0);
}
fn dummy_shutdown<'a>(
grpc: &'a tokio::task::JoinHandle<()>,
http: &'a tokio::task::JoinHandle<()>,
) -> ShutdownTargets<'a> {
ShutdownTargets {
energy: EnergyScraperHandles {
alumet: None,
scaphandre: None,
kepler: None,
redfish: None,
cloud: None,
emaps: None,
},
listeners: ListenerHandles {
grpc,
http,
json_socket: None,
},
}
}
fn test_loop_cfg() -> EventLoopConfig {
EventLoopConfig {
green_enabled: true,
sampling_rate: 1.0,
evict_ms: 60_000,
confidence: Confidence::DaemonStaging,
analysis_queue_capacity: 1024,
waste_sticky_ttl_ms: 0,
}
}
#[tokio::test]
async fn fail_loud_returns_error_when_worker_dies() {
let metrics = MetricsState::new();
let base = Arc::new(empty_carbon_ctx());
let sources = no_scrapers(&base);
let window = test_window();
let (_tx, mut rx) = mpsc::channel::<Vec<SpanEvent>>(16);
let (work_tx, _work_rx) = mpsc::channel::<AnalysisBatch>(4);
let worker = tokio::spawn(async {});
let grpc = tokio::spawn(std::future::pending::<()>());
let http = tokio::spawn(std::future::pending::<()>());
let result = drive_event_loop(
&mut rx,
&window,
&metrics,
&sources,
dummy_shutdown(&grpc, &http),
test_loop_cfg(),
work_tx,
worker,
std::future::pending::<()>(), )
.await;
assert!(matches!(
result,
Err(crate::DaemonError::AnalysisWorkerStopped)
));
}
#[tokio::test]
async fn graceful_shutdown_drains_window_and_returns_ok() {
let metrics = Arc::new(MetricsState::new());
let store = Arc::new(findings_store::FindingsStore::new(100));
let cell = fresh_green_cell();
let base = Arc::new(empty_carbon_ctx());
let sources = no_scrapers(&base);
let window = test_window();
{
let mut w = window.lock().await;
for id in ["w1", "w2", "w3"] {
w.push(make_normalized(id, "SELECT 1"), current_time_ms());
}
}
let (_tx, mut rx) = mpsc::channel::<Vec<SpanEvent>>(16);
let (work_tx, work_rx) = mpsc::channel::<AnalysisBatch>(4);
let worker = tokio::spawn(run_analysis_worker(
work_rx,
test_worker_ctx(&metrics, &store, &cell),
));
let grpc = tokio::spawn(std::future::pending::<()>());
let http = tokio::spawn(std::future::pending::<()>());
let (sd_tx, sd_rx) = tokio::sync::oneshot::channel::<()>();
sd_tx.send(()).expect("receiver alive");
let shutdown_fut = async move {
let _ = sd_rx.await;
};
let result = drive_event_loop(
&mut rx,
&window,
&metrics,
&sources,
dummy_shutdown(&grpc, &http),
test_loop_cfg(),
work_tx,
worker,
shutdown_fut,
)
.await;
assert!(result.is_ok());
assert!((metrics.traces_analyzed_total.get() - 3.0).abs() < f64::EPSILON);
}
}