use crate::calibration::{
CalibrationManifest, ConfigurationFingerprint, AttestationService, AttestationSignature,
Ed25519KeyPair, SloSystem, SloStatus, ECEAlert, RegressionDetector, RegressionSeverity,
};
use anyhow::{Context, Result, bail};
use chrono::{DateTime, Utc, Duration};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, RwLock};
use tokio::fs;
use tracing::{info, warn, debug};
use uuid::Uuid;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum FingerprintStatus {
Validating {
start_time: DateTime<Utc>,
validation_duration: Duration,
},
Green {
validated_at: DateTime<Utc>,
stability_metrics: StabilityMetrics,
},
Red {
failed_at: DateTime<Utc>,
failure_reasons: Vec<ValidationFailure>,
},
Deprecated {
deprecated_at: DateTime<Utc>,
replacement_fingerprint: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct StabilityMetrics {
pub mean_ece: f64,
pub max_ece: f64,
pub ece_std: f64,
pub mean_drift_rate: f64,
pub max_drift: f64,
pub slo_breaches: u32,
pub alert_count: u32,
pub regression_count: u32,
pub stability_score: f64,
pub validation_period: (DateTime<Utc>, DateTime<Utc>),
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum ValidationFailure {
EceThresholdExceeded { observed: f64, threshold: f64 },
DriftRateExceeded { observed: f64, threshold: f64 },
SloBreaches { count: u32, threshold: u32 },
RegressionDetected { severity: RegressionSeverity },
InsufficientStability { duration: Duration, required: Duration },
AlertsTriggered { count: u32, threshold: u32 },
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PublishedFingerprint {
pub fingerprint_id: String,
pub status: FingerprintStatus,
pub manifest: CalibrationManifest,
pub attestation: AttestationSignature,
pub publication_metadata: PublicationMetadata,
pub release_verification: ReleaseVerification,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PublicationMetadata {
pub published_at: DateTime<Utc>,
pub publisher: String,
pub version: String,
pub distribution_channels: Vec<String>,
pub download_urls: HashMap<String, String>,
pub verification_instructions: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ReleaseVerification {
pub artifact_checksums: HashMap<String, String>,
pub verification_script_hash: String,
pub gpg_signature: Option<String>,
pub certificate_chain: Vec<String>,
pub verified_at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct FingerprintPublisherConfig {
pub validation_duration: Duration,
pub ece_threshold: f64,
pub drift_threshold: f64,
pub max_slo_breaches: u32,
pub max_alerts: u32,
pub repository_config: RepositoryConfig,
pub attestation_config: PublisherAttestationConfig,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct RepositoryConfig {
pub base_url: String,
pub auth_token: Option<String>,
pub publication_path: String,
pub repository_type: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PublisherAttestationConfig {
pub signing_key_id: String,
pub certificate_chain_path: PathBuf,
pub hsm_config: Option<HsmConfig>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HsmConfig {
pub slot_id: u32,
pub token_label: String,
pub pin: String,
}
pub struct FingerprintPublisher {
config: FingerprintPublisherConfig,
attestation_service: Arc<AttestationService>,
slo_system: Arc<SloSystem>,
regression_detector: Arc<RegressionDetector>,
published_fingerprints: Arc<RwLock<HashMap<String, PublishedFingerprint>>>,
validation_sessions: Arc<RwLock<HashMap<String, ValidationSession>>>,
}
#[derive(Debug, Clone)]
struct ValidationSession {
session_id: String,
fingerprint_id: String,
manifest: CalibrationManifest,
started_at: DateTime<Utc>,
required_duration: Duration,
metrics: Vec<ValidationMetric>,
status: ValidationSessionStatus,
}
#[derive(Debug, Clone)]
enum ValidationSessionStatus {
Active,
Completed { result: ValidationResult },
Failed { failures: Vec<ValidationFailure> },
}
#[derive(Debug, Clone)]
struct ValidationMetric {
timestamp: DateTime<Utc>,
ece: f64,
drift_rate: f64,
slo_status: SloStatus,
active_alerts: Vec<ECEAlert>,
}
#[derive(Debug, Clone)]
enum ValidationResult {
Passed { metrics: StabilityMetrics },
Failed { failures: Vec<ValidationFailure> },
}
impl Default for FingerprintPublisherConfig {
fn default() -> Self {
Self {
validation_duration: Duration::hours(24),
ece_threshold: 0.015,
drift_threshold: 0.001,
max_slo_breaches: 0,
max_alerts: 5,
repository_config: RepositoryConfig {
base_url: "https://releases.calibration.ai".to_string(),
auth_token: None,
publication_path: "fingerprints/green/".to_string(),
repository_type: "https".to_string(),
},
attestation_config: PublisherAttestationConfig {
signing_key_id: "calibration-publisher-2024".to_string(),
certificate_chain_path: PathBuf::from("certs/publisher-chain.pem"),
hsm_config: None,
},
}
}
}
impl FingerprintPublisher {
pub fn new(
config: FingerprintPublisherConfig,
attestation_service: Arc<AttestationService>,
slo_system: Arc<SloSystem>,
regression_detector: Arc<RegressionDetector>,
) -> Self {
Self {
config,
attestation_service,
slo_system,
regression_detector,
published_fingerprints: Arc::new(RwLock::new(HashMap::new())),
validation_sessions: Arc::new(RwLock::new(HashMap::new())),
}
}
pub async fn start_validation_session(
&self,
manifest: CalibrationManifest,
) -> Result<String> {
let session_id = Uuid::new_v4().to_string();
let fingerprint_id = manifest.config_fingerprint.config_hash.clone();
let session = ValidationSession {
session_id: session_id.clone(),
fingerprint_id: fingerprint_id.clone(),
manifest,
started_at: Utc::now(),
required_duration: self.config.validation_duration,
metrics: Vec::new(),
status: ValidationSessionStatus::Active,
};
{
let mut sessions = self.validation_sessions.write()
.map_err(|e| anyhow::anyhow!("Lock poisoned: {}", e))?;
sessions.insert(session_id.clone(), session);
}
info!(
session_id = %session_id,
fingerprint_id = %fingerprint_id,
duration_hours = %self.config.validation_duration.num_hours(),
"Started 24-hour validation session for fingerprint"
);
Ok(session_id)
}
pub async fn collect_validation_metrics(
&self,
session_id: &str,
ece: f64,
drift_rate: f64,
) -> Result<()> {
let slo_status = self.slo_system.get_current_status().await?;
let active_alerts = self.slo_system.get_active_alerts().await?;
let metric = ValidationMetric {
timestamp: Utc::now(),
ece,
drift_rate,
slo_status,
active_alerts,
};
{
let mut sessions = self.validation_sessions.write()
.map_err(|e| anyhow::anyhow!("Lock poisoned: {}", e))?;
if let Some(session) = sessions.get_mut(session_id) {
session.metrics.push(metric);
debug!(
session_id = %session_id,
ece = %ece,
drift_rate = %drift_rate,
"Collected validation metric"
);
} else {
bail!("Validation session not found: {}", session_id);
}
}
Ok(())
}
pub async fn check_validation_completion(
&self,
session_id: &str,
) -> Result<Option<ValidationResult>> {
let mut sessions = self.validation_sessions.write()
.map_err(|e| anyhow::anyhow!("Lock poisoned: {}", e))?;
if let Some(session) = sessions.get_mut(session_id) {
let elapsed = Utc::now().signed_duration_since(session.started_at);
if elapsed >= session.required_duration {
let result = self.assess_validation_results(&session.metrics).await?;
session.status = ValidationSessionStatus::Completed {
result: result.clone(),
};
info!(
session_id = %session_id,
result = ?result,
"Validation session completed"
);
Ok(Some(result))
} else {
Ok(None)
}
} else {
bail!("Validation session not found: {}", session_id);
}
}
async fn assess_validation_results(
&self,
metrics: &[ValidationMetric],
) -> Result<ValidationResult> {
if metrics.is_empty() {
return Ok(ValidationResult::Failed {
failures: vec![ValidationFailure::InsufficientStability {
duration: Duration::zero(),
required: self.config.validation_duration,
}],
});
}
let mut failures = Vec::new();
let eces: Vec<f64> = metrics.iter().map(|m| m.ece).collect();
let drift_rates: Vec<f64> = metrics.iter().map(|m| m.drift_rate).collect();
let mean_ece = eces.iter().sum::<f64>() / eces.len() as f64;
let max_ece = eces.iter().fold(0.0, |a, &b| a.max(b));
let ece_variance = eces.iter()
.map(|&x| (x - mean_ece).powi(2))
.sum::<f64>() / eces.len() as f64;
let ece_std = ece_variance.sqrt();
let mean_drift_rate = drift_rates.iter().sum::<f64>() / drift_rates.len() as f64;
let max_drift = drift_rates.iter().fold(0.0, |a, &b| a.max(b));
let slo_breaches = metrics.iter()
.filter(|m| !matches!(m.slo_status, SloStatus::Healthy))
.count() as u32;
let total_alerts = metrics.iter()
.map(|m| m.active_alerts.len() as u32)
.sum::<u32>();
let regression_count = self.regression_detector.get_incident_count_since(
metrics[0].timestamp
).await.unwrap_or(0);
if mean_ece > self.config.ece_threshold {
failures.push(ValidationFailure::EceThresholdExceeded {
observed: mean_ece,
threshold: self.config.ece_threshold,
});
}
if mean_drift_rate > self.config.drift_threshold {
failures.push(ValidationFailure::DriftRateExceeded {
observed: mean_drift_rate,
threshold: self.config.drift_threshold,
});
}
if slo_breaches > self.config.max_slo_breaches {
failures.push(ValidationFailure::SloBreaches {
count: slo_breaches,
threshold: self.config.max_slo_breaches,
});
}
if total_alerts > self.config.max_alerts {
failures.push(ValidationFailure::AlertsTriggered {
count: total_alerts,
threshold: self.config.max_alerts,
});
}
if regression_count > 0 {
failures.push(ValidationFailure::RegressionDetected {
severity: RegressionSeverity::High, });
}
if failures.is_empty() {
let ece_score = (self.config.ece_threshold - mean_ece) / self.config.ece_threshold;
let drift_score = (self.config.drift_threshold - mean_drift_rate) / self.config.drift_threshold;
let alert_score = if total_alerts == 0 { 1.0 } else {
1.0 - (total_alerts as f64 / self.config.max_alerts as f64)
};
let stability_score = (ece_score + drift_score + alert_score) / 3.0;
let stability_metrics = StabilityMetrics {
mean_ece,
max_ece,
ece_std,
mean_drift_rate,
max_drift,
slo_breaches,
alert_count: total_alerts,
regression_count,
stability_score: stability_score.max(0.0).min(1.0),
validation_period: (metrics[0].timestamp, metrics.last().unwrap().timestamp),
};
Ok(ValidationResult::Passed { metrics: stability_metrics })
} else {
Ok(ValidationResult::Failed { failures })
}
}
pub async fn publish_green_fingerprint(
&self,
session_id: &str,
) -> Result<PublishedFingerprint> {
let session = {
let sessions = self.validation_sessions.read()
.map_err(|e| anyhow::anyhow!("Lock poisoned: {}", e))?;
sessions.get(session_id).cloned()
.ok_or_else(|| anyhow::anyhow!("Session not found: {}", session_id))?
};
match session.status {
ValidationSessionStatus::Completed { result: ValidationResult::Passed { metrics } } => {
let status = FingerprintStatus::Green {
validated_at: Utc::now(),
stability_metrics: metrics,
};
let attestation = self.attestation_service.create_manifest_attestation(
&session.manifest,
&self.config.attestation_config.signing_key_id,
).await.context("Failed to create attestation")?;
let publication_metadata = PublicationMetadata {
published_at: Utc::now(),
publisher: "Calibration Publisher Service v1.0".to_string(),
version: session.manifest.manifest_version.clone(),
distribution_channels: vec![
"https".to_string(),
"ipfs".to_string(),
],
download_urls: HashMap::from([
("https".to_string(), format!("{}/{}.json",
self.config.repository_config.base_url,
session.fingerprint_id)),
("checksum".to_string(), format!("{}/{}.sha256",
self.config.repository_config.base_url,
session.fingerprint_id)),
]),
verification_instructions: "Verify using: lens-verify --fingerprint {fingerprint_id}".to_string(),
};
let release_verification = self.create_release_verification(
&session.manifest,
&attestation,
).await?;
let published_fingerprint = PublishedFingerprint {
fingerprint_id: session.fingerprint_id.clone(),
status,
manifest: session.manifest.clone(),
attestation,
publication_metadata,
release_verification,
};
{
let mut published = self.published_fingerprints.write()
.map_err(|e| anyhow::anyhow!("Lock poisoned: {}", e))?;
published.insert(session.fingerprint_id.clone(), published_fingerprint.clone());
}
self.publish_to_repository(&published_fingerprint).await?;
info!(
fingerprint_id = %session.fingerprint_id,
stability_score = %published_fingerprint.status.stability_score(),
"Successfully published green fingerprint"
);
Ok(published_fingerprint)
}
_ => bail!("Session is not ready for publishing: validation not passed"),
}
}
async fn create_release_verification(
&self,
manifest: &CalibrationManifest,
attestation: &AttestationSignature,
) -> Result<ReleaseVerification> {
let mut artifact_checksums = HashMap::new();
let manifest_bytes = serde_json::to_vec(manifest)?;
let manifest_hash = Sha256::digest(&manifest_bytes);
artifact_checksums.insert(
"manifest.json".to_string(),
hex::encode(manifest_hash),
);
let attestation_bytes = serde_json::to_vec(attestation)?;
let attestation_hash = Sha256::digest(&attestation_bytes);
artifact_checksums.insert(
"attestation.json".to_string(),
hex::encode(attestation_hash),
);
let verification_script = r#"#!/bin/bash
# Verification script for green fingerprints
echo "Verifying fingerprint integrity..."
"#;
let script_hash = Sha256::digest(verification_script.as_bytes());
Ok(ReleaseVerification {
artifact_checksums,
verification_script_hash: hex::encode(script_hash),
gpg_signature: None, certificate_chain: vec![], verified_at: Utc::now(),
})
}
async fn publish_to_repository(
&self,
published_fingerprint: &PublishedFingerprint,
) -> Result<()> {
let json_content = serde_json::to_string_pretty(published_fingerprint)?;
let filename = format!("{}.json", published_fingerprint.fingerprint_id);
let filepath = PathBuf::from(&self.config.repository_config.publication_path)
.join(&filename);
if let Some(parent) = filepath.parent() {
fs::create_dir_all(parent).await?;
}
fs::write(&filepath, json_content).await
.context("Failed to write published fingerprint")?;
info!(
fingerprint_id = %published_fingerprint.fingerprint_id,
path = %filepath.display(),
"Published fingerprint to repository"
);
Ok(())
}
pub fn get_published_fingerprint(&self, fingerprint_id: &str) -> Result<Option<PublishedFingerprint>> {
let published = self.published_fingerprints.read()
.map_err(|e| anyhow::anyhow!("Lock poisoned: {}", e))?;
Ok(published.get(fingerprint_id).cloned())
}
pub fn list_published_fingerprints(&self) -> Result<Vec<PublishedFingerprint>> {
let published = self.published_fingerprints.read()
.map_err(|e| anyhow::anyhow!("Lock poisoned: {}", e))?;
Ok(published.values().cloned().collect())
}
pub async fn verify_published_fingerprint(
&self,
fingerprint_id: &str,
) -> Result<bool> {
if let Some(published) = self.get_published_fingerprint(fingerprint_id)? {
let verification_result = self.attestation_service.verify_manifest_attestation(
&published.manifest,
&published.attestation,
).await?;
Ok(verification_result.is_valid)
} else {
bail!("Published fingerprint not found: {}", fingerprint_id);
}
}
}
impl FingerprintStatus {
pub fn stability_score(&self) -> f64 {
match self {
FingerprintStatus::Green { stability_metrics, .. } => stability_metrics.stability_score,
_ => 0.0,
}
}
pub fn is_green(&self) -> bool {
matches!(self, FingerprintStatus::Green { .. })
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
#[tokio::test]
async fn test_fingerprint_status_creation() {
let status = FingerprintStatus::Validating {
start_time: Utc::now(),
validation_duration: Duration::hours(24),
};
assert!(!status.is_green());
assert_eq!(status.stability_score(), 0.0);
}
#[tokio::test]
async fn test_stability_metrics() {
let metrics = StabilityMetrics {
mean_ece: 0.010,
max_ece: 0.012,
ece_std: 0.001,
mean_drift_rate: 0.0005,
max_drift: 0.0008,
slo_breaches: 0,
alert_count: 2,
regression_count: 0,
stability_score: 0.95,
validation_period: (Utc::now() - Duration::hours(24), Utc::now()),
};
assert!(metrics.stability_score > 0.9);
assert!(metrics.mean_ece < 0.015);
}
}