#![forbid(unsafe_code)]
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use xbp_delivery::{DeliveryEvent, DeliveryEventKind};
#[derive(Debug, Clone, Copy, PartialEq, PartialOrd, Serialize)]
pub struct Hours(pub f64);
#[derive(Debug, Clone, Copy, PartialEq, PartialOrd, Serialize)]
pub struct Days(pub f64);
#[derive(Debug, Clone, Copy, PartialEq, PartialOrd, Serialize)]
pub struct RatePerDay(pub f64);
#[derive(Debug, Clone, Copy, PartialEq, PartialOrd, Serialize)]
pub struct Probability(pub f64);
#[derive(Debug, Clone, Copy, PartialEq, PartialOrd, Serialize)]
pub struct Coverage(pub f64);
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AnalyticsError {
EmptySample,
InvalidProbability,
InvalidWindow,
NegativeDuration,
NonFiniteValue,
InvalidEvidenceCoverage,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum AnalyticsSubjectKind {
Repository,
Service,
Project,
Initiative,
Portfolio,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct AnalyticsSubject {
pub kind: AnalyticsSubjectKind,
pub id: String,
}
impl AnalyticsSubject {
pub fn new(kind: AnalyticsSubjectKind, id: impl Into<String>) -> Self {
Self {
kind,
id: id.into(),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum TrendDirection {
Improving,
Stable,
Declining,
Volatile,
StructuralBreak,
Unknown,
}
#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
pub struct DoraTrend {
pub current: Option<f64>,
pub previous: Option<f64>,
pub delta_absolute: Option<f64>,
pub delta_relative: Option<f64>,
pub trend: TrendDirection,
pub confidence: f64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DoraComparison {
pub subject: AnalyticsSubject,
pub deployment_frequency: DoraTrend,
pub lead_time: DoraTrend,
pub change_failure_rate: DoraTrend,
pub recovery_time: DoraTrend,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum LittleLawStatus {
Consistent,
Divergent,
Insufficient,
}
#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
pub struct LittleLawDiagnostic {
pub wip: usize,
pub throughput_per_day: f64,
pub cycle_time_hours: f64,
pub expected_wip: f64,
pub relative_error: f64,
pub status: LittleLawStatus,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum WipBand {
Normal,
Elevated,
Overloaded,
Unknown,
}
pub fn classify_wip_band(wip: usize, throughput_per_day: f64, cycle_time_hours: f64) -> WipBand {
if throughput_per_day <= 0.0 || cycle_time_hours <= 0.0 {
return WipBand::Unknown;
}
let expected = throughput_per_day * cycle_time_hours / 24.0;
if expected <= 0.0 {
WipBand::Unknown
} else {
let ratio = wip as f64 / expected;
if ratio <= 1.25 {
WipBand::Normal
} else if ratio <= 2.0 {
WipBand::Elevated
} else {
WipBand::Overloaded
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
pub struct WipAgeAssessment {
pub age_hours: f64,
pub historical_p85_hours: Option<f64>,
pub percentile: Option<f64>,
pub aging: bool,
}
pub fn assess_wip_age(age_hours: f64, historical_cycle_times_hours: &[f64]) -> WipAgeAssessment {
let mut sorted = historical_cycle_times_hours
.iter()
.copied()
.filter(|value| value.is_finite() && *value >= 0.0)
.collect::<Vec<_>>();
sorted.sort_by(f64::total_cmp);
let historical_p85_hours = if sorted.is_empty() {
None
} else {
let index = ((sorted.len() as f64 * 0.85).ceil() as usize)
.saturating_sub(1)
.min(sorted.len() - 1);
Some(sorted[index])
};
let percentile = if sorted.is_empty() {
None
} else {
Some(
sorted.iter().filter(|value| **value <= age_hours).count() as f64 / sorted.len() as f64,
)
};
WipAgeAssessment {
age_hours,
historical_p85_hours,
percentile,
aging: historical_p85_hours.is_some_and(|value| age_hours > value),
}
}
pub fn little_law_diagnostic(
wip: usize,
throughput_per_day: f64,
cycle_time_hours: f64,
) -> LittleLawDiagnostic {
let expected_wip = throughput_per_day * cycle_time_hours / 24.0;
let relative_error = if expected_wip == 0.0 {
if wip == 0 {
0.0
} else {
1.0
}
} else {
(wip as f64 - expected_wip).abs() / expected_wip.abs()
};
LittleLawDiagnostic {
wip,
throughput_per_day,
cycle_time_hours,
expected_wip,
relative_error,
status: if throughput_per_day <= 0.0 || cycle_time_hours <= 0.0 {
LittleLawStatus::Insufficient
} else if relative_error <= 0.25 {
LittleLawStatus::Consistent
} else {
LittleLawStatus::Divergent
},
}
}
pub fn dora_comparison(current: &DoraMetrics, previous: &DoraMetrics) -> DoraComparison {
let subject = current
.subject
.clone()
.unwrap_or_else(|| AnalyticsSubject::new(AnalyticsSubjectKind::Project, "unknown"));
DoraComparison {
subject,
deployment_frequency: dora_trend(
current.deployment_frequency_per_week,
previous.deployment_frequency_per_week,
false,
),
lead_time: dora_trend(
current.lead_time_for_changes_hours,
previous.lead_time_for_changes_hours,
true,
),
change_failure_rate: dora_trend(
current.change_failure_rate,
previous.change_failure_rate,
true,
),
recovery_time: dora_trend(
current.failed_deployment_recovery_hours,
previous.failed_deployment_recovery_hours,
true,
),
}
}
fn dora_trend(current: Option<f64>, previous: Option<f64>, lower_is_better: bool) -> DoraTrend {
let delta_absolute = current
.zip(previous)
.map(|(current, previous)| current - previous);
let delta_relative = current
.zip(previous)
.filter(|(_, previous)| *previous != 0.0)
.map(|(current, previous)| (current - previous) / previous.abs());
let trend = match delta_relative {
None => TrendDirection::Unknown,
Some(delta) if delta.abs() < 0.10 => TrendDirection::Stable,
Some(delta) if (delta < 0.0) == lower_is_better => TrendDirection::Improving,
Some(_) => TrendDirection::Declining,
};
DoraTrend {
current,
previous,
delta_absolute,
delta_relative,
trend,
confidence: if current.is_some() && previous.is_some() {
1.0
} else {
0.0
},
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct TimeWindow {
pub start: DateTime<Utc>,
pub end: DateTime<Utc>,
}
impl TimeWindow {
pub fn new(start: DateTime<Utc>, end: DateTime<Utc>) -> Self {
Self { start, end }
}
pub fn days(&self) -> Result<f64, AnalyticsError> {
let seconds = (self.end - self.start).num_seconds();
if seconds <= 0 {
return Err(AnalyticsError::InvalidWindow);
}
Ok(seconds as f64 / 86_400.0)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
pub enum PercentileMethod {
NearestRank,
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct DistributionSummary {
pub sample_count: usize,
pub method: PercentileMethod,
pub min: Option<f64>,
pub max: Option<f64>,
pub mean: Option<f64>,
pub median: Option<f64>,
pub p50: Option<f64>,
pub p75: Option<f64>,
pub p85: Option<f64>,
pub p90: Option<f64>,
pub p95: Option<f64>,
pub p99: Option<f64>,
pub stddev: Option<f64>,
pub mad: Option<f64>,
pub iqr: Option<f64>,
}
impl DistributionSummary {
pub fn from_samples(values: &[f64]) -> Result<Self, AnalyticsError> {
validate_values(values)?;
if values.is_empty() {
return Ok(Self {
sample_count: 0,
method: PercentileMethod::NearestRank,
min: None,
max: None,
mean: None,
median: None,
p50: None,
p75: None,
p85: None,
p90: None,
p95: None,
p99: None,
stddev: None,
mad: None,
iqr: None,
});
}
Ok(Self {
sample_count: values.len(),
method: PercentileMethod::NearestRank,
min: Some(values.iter().copied().fold(f64::INFINITY, f64::min)),
max: Some(values.iter().copied().fold(f64::NEG_INFINITY, f64::max)),
mean: Some(mean(values)?),
median: Some(percentile(values, 0.5)?),
p50: Some(percentile(values, 0.5)?),
p75: Some(percentile(values, 0.75)?),
p85: Some(percentile(values, 0.85)?),
p90: Some(percentile(values, 0.90)?),
p95: Some(percentile(values, 0.95)?),
p99: Some(percentile(values, 0.99)?),
stddev: Some(variance(values)?.sqrt()),
mad: Some(mad(values)?),
iqr: Some(percentile(values, 0.75)? - percentile(values, 0.25)?),
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum CoverageAvailability {
Available,
#[default]
Unavailable,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct EvidenceCoverage {
pub eligible: usize,
pub observed: usize,
pub sufficient: usize,
pub ratio: f64,
#[serde(default)]
pub availability: CoverageAvailability,
}
impl EvidenceCoverage {
pub fn new(eligible: usize, observed: usize, sufficient: usize) -> Self {
Self::try_new(eligible, observed, sufficient)
.expect("evidence coverage denominators must be ordered")
}
pub fn try_new(
eligible: usize,
observed: usize,
sufficient: usize,
) -> Result<Self, AnalyticsError> {
if observed > eligible || sufficient > observed {
return Err(AnalyticsError::InvalidEvidenceCoverage);
}
Ok(Self {
eligible,
observed,
sufficient,
ratio: if eligible == 0 {
0.0
} else {
sufficient as f64 / eligible as f64
},
availability: if eligible > 0 && sufficient > 0 {
CoverageAvailability::Available
} else {
CoverageAvailability::Unavailable
},
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum DataQualityWarning {
InsufficientHistory,
IncompleteContributorBinding,
MissingDeploymentBindings,
StaleProviderData,
PartialStateHistory,
LowDeploymentCorrelationCoverage,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DataQuality {
pub coverage: EvidenceCoverage,
pub freshness_seconds: Option<i64>,
pub effective_sample_size: f64,
pub missing_fields: Vec<String>,
pub warnings: Vec<DataQualityWarning>,
}
impl DataQuality {
pub fn new(
coverage: EvidenceCoverage,
freshness_seconds: Option<i64>,
effective_sample_size: f64,
missing_fields: Vec<String>,
warnings: Vec<DataQualityWarning>,
) -> Self {
Self {
coverage,
freshness_seconds,
effective_sample_size,
missing_fields,
warnings,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct FlowSamples<'a> {
pub lead_times_hours: &'a [f64],
pub cycle_times_hours: &'a [f64],
pub active_times_hours: &'a [f64],
pub blocked_times_hours: &'a [f64],
pub review_wait_times_hours: &'a [f64],
pub deployment_wait_times_hours: &'a [f64],
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct FlowMetrics {
pub window: TimeWindow,
pub completed_count: usize,
pub eligible_completed_count: usize,
pub observed_completed_count: usize,
pub valid_completed_count: usize,
pub arrival_count: usize,
pub wip: usize,
pub blocked_wip: usize,
pub throughput_per_day: f64,
pub arrival_rate_per_day: f64,
pub queue_growth_per_day: f64,
pub lead_time_hours: DistributionSummary,
pub cycle_time_hours: DistributionSummary,
pub active_time_hours: DistributionSummary,
pub blocked_time_hours: DistributionSummary,
pub review_wait_hours: DistributionSummary,
pub deployment_wait_hours: DistributionSummary,
pub flow_efficiency: Option<f64>,
pub sample_count: usize,
pub coverage: EvidenceCoverage,
}
impl FlowMetrics {
pub const SCHEMA_V1: &'static str = "xbp.analytics.flow/v1";
pub const SCHEMA_V2: &'static str = "xbp.analytics.flow/v2";
pub fn from_samples(
window: TimeWindow,
completed_count: usize,
arrival_count: usize,
wip: usize,
samples: FlowSamples<'_>,
) -> Result<Self, AnalyticsError> {
Self::from_samples_with_coverage(
window,
completed_count,
completed_count,
completed_count,
arrival_count,
wip,
0,
samples,
)
}
#[allow(clippy::too_many_arguments)]
pub fn from_samples_with_coverage(
window: TimeWindow,
eligible_completed_count: usize,
observed_completed_count: usize,
valid_completed_count: usize,
arrival_count: usize,
wip: usize,
blocked_wip: usize,
samples: FlowSamples<'_>,
) -> Result<Self, AnalyticsError> {
reject_negative_durations(&[
samples.lead_times_hours,
samples.cycle_times_hours,
samples.active_times_hours,
samples.blocked_times_hours,
samples.review_wait_times_hours,
samples.deployment_wait_times_hours,
])?;
let window_days = window.days()?;
let cycle = DistributionSummary::from_samples(samples.cycle_times_hours)?;
let active = DistributionSummary::from_samples(samples.active_times_hours)?;
let flow_efficiency = match (active.mean, cycle.mean) {
(Some(active), Some(cycle)) if cycle > 0.0 => Some(active / cycle),
_ => None,
};
let coverage = EvidenceCoverage::try_new(
eligible_completed_count,
observed_completed_count,
valid_completed_count,
)?;
Ok(Self {
window,
completed_count: valid_completed_count,
eligible_completed_count,
observed_completed_count,
valid_completed_count,
arrival_count,
wip,
blocked_wip,
throughput_per_day: valid_completed_count as f64 / window_days,
arrival_rate_per_day: arrival_count as f64 / window_days,
queue_growth_per_day: (arrival_count as f64 - valid_completed_count as f64)
/ window_days,
lead_time_hours: DistributionSummary::from_samples(samples.lead_times_hours)?,
cycle_time_hours: cycle,
active_time_hours: active,
blocked_time_hours: DistributionSummary::from_samples(samples.blocked_times_hours)?,
review_wait_hours: DistributionSummary::from_samples(samples.review_wait_times_hours)?,
deployment_wait_hours: DistributionSummary::from_samples(
samples.deployment_wait_times_hours,
)?,
flow_efficiency,
sample_count: samples.lead_times_hours.len(),
coverage,
})
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct FlowMetricsV1 {
pub throughput: f64,
pub arrival_rate: f64,
pub wip: usize,
pub queue_growth: f64,
pub lead_time: Option<f64>,
pub cycle_time: Option<f64>,
pub active_time: Option<f64>,
pub blocked_time: Option<f64>,
pub review_wait: Option<f64>,
pub deployment_wait: Option<f64>,
pub flow_efficiency: Option<f64>,
pub sample_count: usize,
}
impl FlowMetrics {
pub fn to_v1(&self) -> FlowMetricsV1 {
FlowMetricsV1 {
throughput: self.throughput_per_day,
arrival_rate: self.arrival_rate_per_day,
wip: self.wip,
queue_growth: self.queue_growth_per_day,
lead_time: self.lead_time_hours.mean,
cycle_time: self.cycle_time_hours.mean,
active_time: self.active_time_hours.mean,
blocked_time: self.blocked_time_hours.mean,
review_wait: self.review_wait_hours.mean,
deployment_wait: self.deployment_wait_hours.mean,
flow_efficiency: self.flow_efficiency,
sample_count: self.sample_count,
}
}
}
pub type FlowProjection = FlowMetrics;
impl FlowMetrics {
pub fn from_events(events: &[DeliveryEvent], window: TimeWindow) -> Result<Self, String> {
use std::collections::{BTreeMap, BTreeSet};
window
.days()
.map_err(|error| format!("invalid flow window: {error:?}"))?;
let mut histories: BTreeMap<String, Vec<&DeliveryEvent>> = BTreeMap::new();
for event in events
.iter()
.filter(|event| event.occurred_at <= window.end)
{
let Some(work_item_id) =
event
.work_item_id
.as_ref()
.map(|id| id.0.clone())
.or_else(|| {
event
.attributes
.get("work_item_id")
.and_then(|value| value.as_str())
.map(str::to_string)
})
else {
continue;
};
histories.entry(work_item_id).or_default().push(event);
}
for history in histories.values_mut() {
history.sort_by_key(|event| (event.occurred_at, event.id.clone()));
}
let completed_ids = histories
.iter()
.filter(|(_, history)| {
history.iter().any(|event| {
event.kind == DeliveryEventKind::WorkItemCompleted
&& event.occurred_at >= window.start
&& event.occurred_at <= window.end
})
})
.map(|(id, _)| id.clone())
.collect::<BTreeSet<_>>();
let arrival_count = histories
.values()
.filter(|history| {
history.iter().any(|event| {
event.kind == DeliveryEventKind::WorkItemCreated
&& event.occurred_at >= window.start
&& event.occurred_at <= window.end
})
})
.count();
let mut lead_times = Vec::new();
let mut cycle_times = Vec::new();
let mut blocked_times = Vec::new();
let mut observed_completed_count = 0;
let mut blocked_wip = 0;
for (work_item_id, history) in &histories {
let created = history
.iter()
.find(|event| event.kind == DeliveryEventKind::WorkItemCreated)
.map(|event| event.occurred_at);
let started = history
.iter()
.find(|event| event.kind == DeliveryEventKind::WorkItemStarted)
.map(|event| event.occurred_at);
let completed = history
.iter()
.find(|event| event.kind == DeliveryEventKind::WorkItemCompleted)
.map(|event| event.occurred_at);
let mut is_blocked = false;
for event in history {
match event.kind {
DeliveryEventKind::WorkItemBlocked if !is_blocked => is_blocked = true,
DeliveryEventKind::WorkItemUnblocked if is_blocked => is_blocked = false,
DeliveryEventKind::WorkItemBlocked => {
return Err(format!(
"duplicate block transition for work item {}",
event
.work_item_id
.as_ref()
.map(|id| id.0.as_str())
.unwrap_or("unknown")
));
}
DeliveryEventKind::WorkItemUnblocked => {
return Err("unmatched unblock transition in delivery evidence".into());
}
_ => {}
}
}
if completed.is_none()
&& created.is_some_and(|occurred_at| occurred_at <= window.end)
&& is_blocked
{
blocked_wip += 1;
}
if !completed_ids.contains(work_item_id) {
continue;
}
let Some(created) = created else {
continue;
};
let Some(started) = started else {
continue;
};
let Some(completed) = completed else {
continue;
};
if started < created || completed < started {
return Err("negative lifecycle duration in delivery evidence".into());
}
let (mut blocked_at, mut blocked_hours) = (None, 0.0);
for event in history {
match event.kind {
DeliveryEventKind::WorkItemBlocked => blocked_at = Some(event.occurred_at),
DeliveryEventKind::WorkItemUnblocked => {
let Some(blocked_at_value) = blocked_at.take() else {
return Err("unmatched unblock transition in delivery evidence".into());
};
let hours =
(event.occurred_at - blocked_at_value).num_seconds() as f64 / 3600.0;
if hours < 0.0 {
return Err("negative blocked duration in delivery evidence".into());
}
blocked_hours += hours;
}
_ => {}
}
}
if blocked_at.is_some() {
return Err("unmatched block transition in delivery evidence".into());
}
lead_times.push((completed - created).num_seconds() as f64 / 3600.0);
cycle_times.push((completed - started).num_seconds() as f64 / 3600.0);
blocked_times.push(blocked_hours);
observed_completed_count += 1;
}
Self::from_samples_with_coverage(
window.clone(),
completed_ids.len(),
observed_completed_count,
observed_completed_count,
arrival_count,
histories
.values()
.filter(|history| {
let created = history
.iter()
.find(|event| event.kind == DeliveryEventKind::WorkItemCreated)
.map(|event| event.occurred_at);
let completed = history
.iter()
.find(|event| event.kind == DeliveryEventKind::WorkItemCompleted)
.map(|event| event.occurred_at);
created.is_some_and(|occurred_at| occurred_at <= window.end)
&& completed.is_none()
})
.count(),
blocked_wip,
FlowSamples {
lead_times_hours: &lead_times,
cycle_times_hours: &cycle_times,
active_times_hours: &[],
blocked_times_hours: &blocked_times,
review_wait_times_hours: &[],
deployment_wait_times_hours: &[],
},
)
.map_err(|error| format!("flow projection failed: {error:?}"))
}
}
pub fn count(values: &[f64]) -> usize {
values.len()
}
pub fn mean(values: &[f64]) -> Result<f64, AnalyticsError> {
validate_values(values)?;
if values.is_empty() {
return Err(AnalyticsError::EmptySample);
}
Ok(values.iter().sum::<f64>() / values.len() as f64)
}
pub fn median(values: &[f64]) -> Result<f64, AnalyticsError> {
percentile(values, 0.5)
}
pub fn percentile(values: &[f64], probability: f64) -> Result<f64, AnalyticsError> {
validate_values(values)?;
if values.is_empty() {
return Err(AnalyticsError::EmptySample);
}
if !(0.0..=1.0).contains(&probability) {
return Err(AnalyticsError::InvalidProbability);
}
let mut sorted = values.to_vec();
sorted.sort_by(f64::total_cmp);
let rank = (probability * sorted.len() as f64).ceil().max(1.0) as usize;
Ok(sorted[rank - 1])
}
pub fn variance(values: &[f64]) -> Result<f64, AnalyticsError> {
let average = mean(values)?;
Ok(values
.iter()
.map(|value| (value - average).powi(2))
.sum::<f64>()
/ values.len() as f64)
}
pub fn standard_deviation(values: &[f64]) -> Result<f64, AnalyticsError> {
Ok(variance(values)?.sqrt())
}
pub fn iqr(values: &[f64]) -> Result<f64, AnalyticsError> {
Ok(percentile(values, 0.75)? - percentile(values, 0.25)?)
}
pub fn mad(values: &[f64]) -> Result<f64, AnalyticsError> {
let center = median(values)?;
let deviations: Vec<f64> = values.iter().map(|value| (value - center).abs()).collect();
median(&deviations)
}
pub fn histogram(values: &[f64], bins: usize) -> Result<Vec<usize>, AnalyticsError> {
validate_values(values)?;
if bins == 0 {
return Ok(Vec::new());
}
if values.is_empty() {
return Ok(vec![0; bins]);
}
let min = values.iter().copied().fold(f64::INFINITY, f64::min);
let max = values.iter().copied().fold(f64::NEG_INFINITY, f64::max);
if min == max {
let mut result = vec![0; bins];
result[0] = values.len();
return Ok(result);
}
let width = (max - min) / bins as f64;
let mut result = vec![0; bins];
for value in values {
let index = (((value - min) / width).floor() as usize).min(bins - 1);
result[index] += 1;
}
Ok(result)
}
pub fn empirical_cdf(values: &[f64], threshold: f64) -> Result<f64, AnalyticsError> {
validate_values(values)?;
if values.is_empty() {
return Err(AnalyticsError::EmptySample);
}
if !threshold.is_finite() {
return Err(AnalyticsError::NonFiniteValue);
}
Ok(values.iter().filter(|value| **value <= threshold).count() as f64 / values.len() as f64)
}
fn validate_values(values: &[f64]) -> Result<(), AnalyticsError> {
if values.iter().any(|value| !value.is_finite()) {
return Err(AnalyticsError::NonFiniteValue);
}
Ok(())
}
fn reject_negative_durations(samples: &[&[f64]]) -> Result<(), AnalyticsError> {
for sample in samples {
validate_values(sample)?;
if sample.iter().any(|value| *value < 0.0) {
return Err(AnalyticsError::NegativeDuration);
}
}
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum MetricAvailability {
Available,
Unavailable,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum MetricEvidenceStatus {
#[default]
NoEvidence,
InsufficientEvidence,
Available,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DoraMetricEvidence {
pub eligible: usize,
pub observed: usize,
pub sufficient: usize,
pub availability: MetricAvailability,
}
impl DoraMetricEvidence {
pub fn from_counts(eligible: usize, observed: usize, sufficient: usize) -> Self {
Self::with_availability(eligible, observed, sufficient, sufficient > 0)
}
pub fn status(&self) -> MetricEvidenceStatus {
if self.eligible == 0 {
MetricEvidenceStatus::NoEvidence
} else if self.sufficient < self.eligible {
MetricEvidenceStatus::InsufficientEvidence
} else {
MetricEvidenceStatus::Available
}
}
fn new(eligible: usize, observed: usize, sufficient: usize) -> Self {
Self::from_counts(eligible, observed, sufficient)
}
fn with_availability(
eligible: usize,
observed: usize,
sufficient: usize,
available: bool,
) -> Self {
Self {
eligible,
observed,
sufficient,
availability: if available {
MetricAvailability::Available
} else {
MetricAvailability::Unavailable
},
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DoraMetricsV1 {
pub deployment_frequency_per_week: f64,
pub lead_time_for_changes_hours: Option<f64>,
pub change_failure_rate: f64,
pub failed_deployment_recovery_hours: Option<f64>,
pub sample_count: usize,
pub window_start: DateTime<Utc>,
pub window_end: DateTime<Utc>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DoraMetrics {
#[serde(default)]
pub subject: Option<AnalyticsSubject>,
pub deployment_frequency_per_week: Option<f64>,
pub lead_time_for_changes_hours: Option<f64>,
pub change_failure_rate: Option<f64>,
pub failed_deployment_recovery_hours: Option<f64>,
pub deployment_frequency_evidence: DoraMetricEvidence,
pub lead_time_evidence: DoraMetricEvidence,
pub change_failure_rate_evidence: DoraMetricEvidence,
pub recovery_evidence: DoraMetricEvidence,
#[serde(default)]
pub deployment_frequency_status: MetricEvidenceStatus,
#[serde(default)]
pub lead_time_status: MetricEvidenceStatus,
#[serde(default)]
pub change_failure_rate_status: MetricEvidenceStatus,
#[serde(default)]
pub recovery_status: MetricEvidenceStatus,
#[serde(skip)]
pub sample_count: usize,
pub window_start: DateTime<Utc>,
pub window_end: DateTime<Utc>,
}
impl DoraMetrics {
pub const SCHEMA_V1: &'static str = "xbp.analytics.dora/v1";
pub const SCHEMA_V2: &'static str = "xbp.analytics.dora/v2";
pub fn to_v1(&self) -> DoraMetricsV1 {
DoraMetricsV1 {
deployment_frequency_per_week: self.deployment_frequency_per_week.unwrap_or(0.0),
lead_time_for_changes_hours: self.lead_time_for_changes_hours,
change_failure_rate: self.change_failure_rate.unwrap_or(0.0),
failed_deployment_recovery_hours: self.failed_deployment_recovery_hours,
sample_count: self.sample_count,
window_start: self.window_start,
window_end: self.window_end,
}
}
pub fn from_events(
events: &[DeliveryEvent],
window_start: DateTime<Utc>,
window_end: DateTime<Utc>,
) -> Self {
let in_window = events
.iter()
.filter(|event| event.occurred_at >= window_start && event.occurred_at <= window_end)
.collect::<Vec<_>>();
let production_deployments = in_window
.iter()
.filter(|event| {
matches!(
event.kind,
DeliveryEventKind::DeploymentSucceeded | DeliveryEventKind::DeploymentFailed
) && event
.attributes
.get("environment")
.and_then(|value| value.as_str())
.is_some_and(|environment| environment.eq_ignore_ascii_case("production"))
})
.collect::<Vec<_>>();
let successful_deployments = production_deployments
.iter()
.filter(|event| {
event.kind == DeliveryEventKind::DeploymentSucceeded
&& event
.attributes
.get("outcome")
.and_then(|value| value.as_str())
.map(|outcome| outcome.eq_ignore_ascii_case("succeeded"))
.unwrap_or(true)
})
.collect::<Vec<_>>();
let failed_deployments = production_deployments
.iter()
.filter(|event| has_failed_deployment_outcome(event))
.collect::<Vec<_>>();
let deployment_count = production_deployments.len();
let commits = events
.iter()
.filter(|event| event.kind == DeliveryEventKind::CommitCreated)
.collect::<Vec<_>>();
let correlated_deployments = successful_deployments
.iter()
.filter(|deployment| {
deployment
.attributes
.get("commit_ids")
.and_then(|value| value.as_array())
.is_some_and(|commit_ids| {
commits.iter().any(|commit| {
commit_ids.iter().any(|commit_id| {
commit_id.as_str() == Some(commit.id.as_str())
|| commit_id.as_str() == Some(commit.source_event_id.as_str())
})
})
})
})
.count();
let lead_times = successful_deployments
.iter()
.filter_map(|deployment| {
let commit_ids = deployment
.attributes
.get("commit_ids")
.and_then(|value| value.as_array())?;
commits
.iter()
.filter(|commit| {
commit_ids.iter().any(|commit_id| {
commit_id.as_str() == Some(commit.id.as_str())
|| commit_id.as_str() == Some(commit.source_event_id.as_str())
})
})
.min_by_key(|commit| commit.occurred_at)
.map(|commit| {
(deployment.occurred_at - commit.occurred_at).num_seconds() as f64 / 3600.0
})
})
.filter(|duration| *duration >= 0.0)
.collect::<Vec<_>>();
let correlated_failures = failed_deployments
.iter()
.filter(|failed| {
failed
.attributes
.get("incident_id")
.and_then(|value| value.as_str())
.is_some_and(|incident_id| {
events.iter().any(|event| {
event.kind == DeliveryEventKind::IncidentOpened
&& (event.id == incident_id || event.source_event_id == incident_id)
})
})
})
.count();
let recovery_times = failed_deployments
.iter()
.filter_map(|failed| {
let incident_id = failed
.attributes
.get("incident_id")
.and_then(|value| value.as_str())?;
let opened = events.iter().find(|event| {
event.kind == DeliveryEventKind::IncidentOpened
&& (event.id == incident_id || event.source_event_id == incident_id)
})?;
events
.iter()
.find(|event| {
event.kind == DeliveryEventKind::IncidentResolved
&& (event.id == incident_id || event.source_event_id == incident_id)
})
.map(|commit| {
(commit.occurred_at - opened.occurred_at).num_seconds() as f64 / 3600.0
})
})
.collect::<Vec<_>>();
let window_days = (window_end - window_start).num_seconds() as f64 / 86_400.0;
Self {
subject: None,
deployment_frequency_per_week: (window_days > 0.0)
.then(|| production_deployments.len() as f64 * 7.0 / window_days)
.filter(|value| value.is_finite()),
lead_time_for_changes_hours: mean(&lead_times).ok(),
change_failure_rate: (deployment_count > 0)
.then(|| failed_deployments.len() as f64 / deployment_count as f64),
failed_deployment_recovery_hours: mean(&recovery_times).ok(),
deployment_frequency_evidence: DoraMetricEvidence::with_availability(
deployment_count,
deployment_count,
deployment_count,
window_days > 0.0,
),
lead_time_evidence: DoraMetricEvidence::new(
successful_deployments.len(),
correlated_deployments,
lead_times.len(),
),
change_failure_rate_evidence: DoraMetricEvidence::new(
deployment_count,
deployment_count,
deployment_count,
),
recovery_evidence: DoraMetricEvidence::new(
failed_deployments.len(),
correlated_failures,
recovery_times.len(),
),
deployment_frequency_status: DoraMetricEvidence::from_counts(
deployment_count,
deployment_count,
deployment_count,
)
.status(),
lead_time_status: DoraMetricEvidence::from_counts(
successful_deployments.len(),
correlated_deployments,
lead_times.len(),
)
.status(),
change_failure_rate_status: DoraMetricEvidence::from_counts(
deployment_count,
deployment_count,
deployment_count,
)
.status(),
recovery_status: DoraMetricEvidence::from_counts(
failed_deployments.len(),
correlated_failures,
recovery_times.len(),
)
.status(),
sample_count: lead_times.len(),
window_start,
window_end,
}
}
pub fn from_events_for_subject(
events: &[DeliveryEvent],
window_start: DateTime<Utc>,
window_end: DateTime<Utc>,
subject: &AnalyticsSubject,
) -> Self {
let scoped = events
.iter()
.filter(|event| subject_matches(event, subject))
.cloned()
.collect::<Vec<_>>();
let mut metrics = Self::from_events(&scoped, window_start, window_end);
metrics.subject = Some(subject.clone());
metrics
}
}
fn subject_matches(event: &DeliveryEvent, subject: &AnalyticsSubject) -> bool {
match subject.kind {
AnalyticsSubjectKind::Project => event
.project_id
.as_ref()
.is_some_and(|value| value.0 == subject.id),
AnalyticsSubjectKind::Repository => event
.repository_id
.as_ref()
.is_some_and(|value| value.0 == subject.id),
AnalyticsSubjectKind::Service
| AnalyticsSubjectKind::Initiative
| AnalyticsSubjectKind::Portfolio => {
event
.attributes
.get(match subject.kind {
AnalyticsSubjectKind::Service => "service_id",
AnalyticsSubjectKind::Initiative => "initiative_id",
AnalyticsSubjectKind::Portfolio => "portfolio_id",
AnalyticsSubjectKind::Project | AnalyticsSubjectKind::Repository => {
unreachable!()
}
})
.and_then(|value| value.as_str())
== Some(subject.id.as_str())
}
}
}
fn has_failed_deployment_outcome(event: &DeliveryEvent) -> bool {
event.kind == DeliveryEventKind::DeploymentFailed
|| event
.attributes
.get("outcome")
.and_then(|value| value.as_str())
.map(|outcome| {
matches!(
outcome.to_ascii_lowercase().as_str(),
"failed" | "rolled_back" | "degraded"
)
})
.unwrap_or(false)
}