use serde::Serialize;
use std::collections::HashMap;
use std::sync::Mutex;
use std::sync::atomic::{AtomicI64, AtomicU64, Ordering};
use std::time::{Duration, Instant};
#[derive(Debug, Serialize, Clone, PartialEq, Default)]
pub enum SinkStatus {
Healthy,
Degraded { reason: String },
Unhealthy { error: String },
#[default]
NotStarted,
}
impl SinkStatus {
pub fn is_operational(&self) -> bool {
match self {
SinkStatus::Healthy => true,
SinkStatus::Degraded { .. } => true,
SinkStatus::Unhealthy { .. } => false,
SinkStatus::NotStarted => false,
}
}
fn is_fully_healthy(&self) -> bool {
self == &SinkStatus::Healthy
}
}
#[derive(Debug, Serialize, Clone)]
pub struct SinkHealth {
pub status: SinkStatus,
pub last_error: Option<String>,
pub consecutive_failures: u32,
}
impl Default for SinkHealth {
fn default() -> Self {
Self {
status: SinkStatus::NotStarted,
last_error: None,
consecutive_failures: 0,
}
}
}
impl SinkHealth {
pub fn healthy() -> Self {
Self {
status: SinkStatus::Healthy,
last_error: None,
consecutive_failures: 0,
}
}
pub fn unhealthy(error: String) -> Self {
Self {
status: SinkStatus::Unhealthy {
error: error.clone(),
},
last_error: Some(error),
consecutive_failures: 1,
}
}
}
#[derive(Debug)]
pub struct Gauge {
value: AtomicI64,
}
impl Gauge {
pub fn new(val: i64) -> Self {
Self {
value: AtomicI64::new(val),
}
}
pub fn set(&self, v: i64) {
self.value.store(v, Ordering::Relaxed);
}
pub fn get(&self) -> i64 {
self.value.load(Ordering::Relaxed)
}
pub fn inc(&self) {
self.value.fetch_add(1, Ordering::Relaxed);
}
pub fn dec(&self) {
self.value.fetch_sub(1, Ordering::Relaxed);
}
}
#[derive(Debug)]
pub struct GaugeF64 {
value: Mutex<f64>,
}
impl GaugeF64 {
pub fn new(val: f64) -> Self {
Self {
value: Mutex::new(val),
}
}
pub fn set(&self, v: f64) {
let mut guard = match self.value.lock() {
Ok(guard) => guard,
Err(poisoned) => {
tracing::warn!("GaugeF64 mutex poisoned, recovering");
poisoned.into_inner()
}
};
*guard = v;
}
pub fn get(&self) -> f64 {
let guard = match self.value.lock() {
Ok(guard) => guard,
Err(poisoned) => {
tracing::warn!("GaugeF64 mutex poisoned, recovering");
poisoned.into_inner()
}
};
*guard
}
}
#[derive(Debug)]
pub struct Histogram {
buckets: Vec<AtomicU64>,
bounds: Vec<u64>, }
impl Histogram {
pub fn new(bounds: Vec<u64>) -> Self {
let mut buckets = Vec::with_capacity(bounds.len() + 1);
for _ in 0..=bounds.len() {
buckets.push(AtomicU64::new(0));
}
Self { buckets, bounds }
}
pub fn record(&self, value: u64) {
let mut index = self.bounds.len();
for (i, &bound) in self.bounds.iter().enumerate() {
if value < bound {
index = i;
break;
}
}
self.buckets[index].fetch_add(1, Ordering::Relaxed);
}
pub fn snapshot(&self) -> Vec<u64> {
self.buckets
.iter()
.map(|b| b.load(Ordering::Relaxed))
.collect()
}
pub fn percentile(&self, p: f64) -> u64 {
let snapshot = self.snapshot();
let total: u64 = snapshot.iter().sum();
if total == 0 {
return 0;
}
let target = (total as f64 * p / 100.0).ceil() as u64;
let mut cumulative: u64 = 0;
for (i, &count) in snapshot.iter().enumerate() {
cumulative += count;
if cumulative >= target {
if i < self.bounds.len() {
return self.bounds[i];
} else {
return self.bounds.last().copied().unwrap_or(0);
}
}
}
self.bounds.last().copied().unwrap_or(0)
}
pub fn p50(&self) -> u64 {
self.percentile(50.0)
}
pub fn p95(&self) -> u64 {
self.percentile(95.0)
}
pub fn p99(&self) -> u64 {
self.percentile(99.0)
}
}
#[derive(Debug, Serialize)]
pub struct MetricsSnapshot {
pub logs_written: u64,
pub logs_dropped: u64,
pub channel_blocked: u64,
pub sink_errors: u64,
pub db_batch_size: i64,
pub db_batch_records_total: u64,
pub avg_latency_us: u64,
pub p50_latency_us: u64,
pub p95_latency_us: u64,
pub p99_latency_us: u64,
pub latency_distribution: Vec<u64>,
pub active_workers: i64,
pub pool_hit_rate: f64,
}
#[derive(Debug, Serialize, Clone)]
pub struct PoolStats {
pub log_record_pool_size: usize,
pub string_buffer_pool_size: usize,
pub hit_rate: f64,
}
#[derive(Debug, Serialize)]
pub struct HealthStatus {
pub overall_status: SinkStatus,
pub sinks: HashMap<String, SinkHealth>,
pub channel_usage: f64,
pub uptime_seconds: u64,
pub metrics: MetricsSnapshot,
pub pool_stats: Option<PoolStats>,
pub encryption_key_valid: bool,
}
#[derive(Debug)]
pub struct Metrics {
pub(crate) logs_written_total: AtomicU64,
pub(crate) logs_dropped_total: AtomicU64,
pub(crate) channel_send_blocked_total: AtomicU64,
pub(crate) sink_errors_total: AtomicU64,
pub(crate) lock_contention_total: AtomicU64,
pub(crate) db_batch_records_total: AtomicU64,
pub(crate) start_time: Instant,
pub(crate) total_latency_us: AtomicU64,
pub(crate) latency_count: AtomicU64,
pub(crate) latency_histogram: Histogram,
pub(crate) active_workers: Gauge,
pub(crate) db_batch_size: Gauge,
pub(crate) pool_hit_rate: GaugeF64,
pub(crate) sink_health: Mutex<HashMap<String, SinkHealth>>,
}
impl Default for Metrics {
fn default() -> Self {
let bounds = vec![1000, 5000, 10000, 50000, 100000, 500000, 1000000];
Self {
logs_written_total: AtomicU64::new(0),
logs_dropped_total: AtomicU64::new(0),
channel_send_blocked_total: AtomicU64::new(0),
sink_errors_total: AtomicU64::new(0),
lock_contention_total: AtomicU64::new(0),
db_batch_records_total: AtomicU64::new(0),
start_time: Instant::now(),
total_latency_us: AtomicU64::new(0),
latency_count: AtomicU64::new(0),
latency_histogram: Histogram::new(bounds),
active_workers: Gauge::new(0),
db_batch_size: Gauge::new(0),
pool_hit_rate: GaugeF64::new(0.0),
sink_health: Mutex::new(HashMap::new()),
}
}
}
impl Metrics {
pub fn new() -> Self {
Self::default()
}
#[inline]
fn audit_access(&self, field: &str) {
tracing::debug!(event = "internal_state_access", field = field,);
}
pub fn logs_written(&self) -> u64 {
self.logs_written_total.load(Ordering::Relaxed)
}
pub fn logs_dropped(&self) -> u64 {
self.logs_dropped_total.load(Ordering::Relaxed)
}
pub fn channel_blocked(&self) -> u64 {
self.channel_send_blocked_total.load(Ordering::Relaxed)
}
pub fn sink_errors(&self) -> u64 {
self.sink_errors_total.load(Ordering::Relaxed)
}
pub fn db_batch_size(&self) -> i64 {
self.db_batch_size.get()
}
pub fn db_batch_records_total(&self) -> u64 {
self.db_batch_records_total.load(Ordering::Relaxed)
}
pub fn active_workers(&self) -> i64 {
self.audit_access("active_workers");
self.active_workers.get()
}
pub fn pool_hit_rate(&self) -> f64 {
self.audit_access("pool_hit_rate");
self.pool_hit_rate.get()
}
pub fn set_pool_hit_rate(&self, rate: f64) {
self.pool_hit_rate.set(rate);
}
pub fn sink_health(&self) -> std::collections::HashMap<String, SinkHealth> {
self.audit_access("sink_health");
match self.sink_health.lock() {
Ok(guard) => guard.clone(),
Err(_) => std::collections::HashMap::new(),
}
}
pub fn uptime(&self) -> Duration {
self.start_time.elapsed()
}
pub fn inc_logs_written(&self) {
self.logs_written_total.fetch_add(1, Ordering::Relaxed);
}
pub fn inc_logs_dropped(&self) {
self.logs_dropped_total.fetch_add(1, Ordering::Relaxed);
}
pub fn inc_channel_blocked(&self) {
self.channel_send_blocked_total
.fetch_add(1, Ordering::Relaxed);
}
pub fn inc_sink_error(&self) {
self.sink_errors_total.fetch_add(1, Ordering::Relaxed);
}
pub fn inc_lock_contention(&self) {
self.lock_contention_total.fetch_add(1, Ordering::Relaxed);
}
pub fn lock_contention(&self) -> u64 {
self.lock_contention_total.load(Ordering::Relaxed)
}
pub fn set_db_batch_size(&self, size: usize) {
self.db_batch_size.set(size as i64);
}
pub fn add_db_batch_records_total(&self, count: usize) {
self.db_batch_records_total
.fetch_add(count as u64, Ordering::Relaxed);
}
pub fn record_latency(&self, duration: Duration) {
let micros = duration.as_micros() as u64;
self.total_latency_us.fetch_add(micros, Ordering::Relaxed);
self.latency_count.fetch_add(1, Ordering::Relaxed);
self.latency_histogram.record(micros);
}
pub fn update_sink_health(&self, name: &str, healthy: bool, error: Option<String>) {
let status = if healthy {
SinkStatus::Healthy
} else {
let error_msg = error
.as_ref()
.unwrap_or(&"Unknown error".to_string())
.clone();
SinkStatus::Unhealthy { error: error_msg }
};
let (new_failures, new_error) = if healthy {
(0, None)
} else {
let current_failures = if let Ok(map) = self.sink_health.lock() {
map.get(name).map(|h| h.consecutive_failures).unwrap_or(0)
} else {
0
};
(current_failures + 1, error)
};
if let Ok(mut map) = self.sink_health.lock() {
let entry = map
.entry(name.to_string())
.or_insert_with(SinkHealth::healthy);
entry.status = status;
entry.consecutive_failures = new_failures;
entry.last_error = new_error;
}
}
pub fn sink_started(&self, name: &str) {
if let Ok(mut map) = self.sink_health.lock() {
let entry = map.entry(name.to_string()).or_insert(SinkHealth::healthy());
entry.status = SinkStatus::Healthy;
entry.consecutive_failures = 0;
entry.last_error = None;
}
}
pub fn sink_degraded(&self, name: &str, reason: String) {
if let Ok(mut map) = self.sink_health.lock() {
let entry = map.entry(name.to_string()).or_insert(SinkHealth::healthy());
entry.status = SinkStatus::Degraded {
reason: reason.clone(),
};
entry.last_error = Some(reason);
}
}
pub fn get_status(&self, channel_len: usize, channel_cap: usize) -> HealthStatus {
let sinks: std::collections::HashMap<String, SinkHealth> = match self.sink_health.lock() {
Ok(guard) => guard.clone(),
Err(_e) => {
eprintln!("Metrics mutex poisoned, using empty data");
std::collections::HashMap::new()
}
};
let overall_status = if sinks.is_empty() {
SinkStatus::NotStarted
} else {
let all_healthy = sinks.values().all(|s| s.status.is_fully_healthy());
let any_unhealthy = sinks
.values()
.any(|s| matches!(s.status, SinkStatus::Unhealthy { .. }));
let any_degraded = sinks
.values()
.any(|s| matches!(s.status, SinkStatus::Degraded { .. }));
if all_healthy {
SinkStatus::Healthy
} else if any_unhealthy {
let errors: Vec<String> = sinks
.values()
.filter_map(|s| {
if let SinkStatus::Unhealthy { error } = &s.status {
Some(error.clone())
} else {
None
}
})
.collect();
SinkStatus::Unhealthy {
error: errors.join("; "),
}
} else if any_degraded {
let reasons: Vec<String> = sinks
.values()
.filter_map(|s| {
if let SinkStatus::Degraded { reason } = &s.status {
Some(reason.clone())
} else {
None
}
})
.collect();
SinkStatus::Degraded {
reason: reasons.join("; "),
}
} else {
SinkStatus::Healthy
}
};
let count = self.latency_count.load(Ordering::Relaxed);
let total = self.total_latency_us.load(Ordering::Relaxed);
let avg_latency = total.checked_div(count).unwrap_or(0);
HealthStatus {
overall_status,
sinks,
channel_usage: if channel_cap > 0 {
channel_len as f64 / channel_cap as f64
} else {
0.0
},
uptime_seconds: self.uptime().as_secs(),
metrics: MetricsSnapshot {
logs_written: self.logs_written_total.load(Ordering::Relaxed),
logs_dropped: self.logs_dropped_total.load(Ordering::Relaxed),
channel_blocked: self.channel_send_blocked_total.load(Ordering::Relaxed),
sink_errors: self.sink_errors_total.load(Ordering::Relaxed),
db_batch_size: self.db_batch_size.get(),
db_batch_records_total: self.db_batch_records_total.load(Ordering::Relaxed),
avg_latency_us: avg_latency,
p50_latency_us: self.latency_histogram.p50(),
p95_latency_us: self.latency_histogram.p95(),
p99_latency_us: self.latency_histogram.p99(),
latency_distribution: self.latency_histogram.snapshot(),
active_workers: self.active_workers.get(),
pool_hit_rate: self.pool_hit_rate.get(),
},
pool_stats: None,
encryption_key_valid: true,
}
}
pub fn export_prometheus(&self) -> String {
let mut s = String::new();
s.push_str("# HELP inklog_logs_written_total Total logs successfully written\n");
s.push_str("# TYPE inklog_logs_written_total counter\n");
s.push_str(&format!(
"inklog_logs_written_total {}\n",
self.logs_written_total.load(Ordering::Relaxed)
));
s.push_str("# HELP inklog_logs_dropped_total Total logs dropped\n");
s.push_str("# TYPE inklog_logs_dropped_total counter\n");
s.push_str(&format!(
"inklog_logs_dropped_total {}\n",
self.logs_dropped_total.load(Ordering::Relaxed)
));
s.push_str("# HELP inklog_channel_blocked_total Total times channel was blocked\n");
s.push_str("# TYPE inklog_channel_blocked_total counter\n");
s.push_str(&format!(
"inklog_channel_blocked_total {}\n",
self.channel_send_blocked_total.load(Ordering::Relaxed)
));
s.push_str("# HELP inklog_sink_errors_total Total sink errors\n");
s.push_str("# TYPE inklog_sink_errors_total counter\n");
s.push_str(&format!(
"inklog_sink_errors_total {}\n",
self.sink_errors_total.load(Ordering::Relaxed)
));
s.push_str("# HELP inklog_db_batch_size Database batch size in last flush\n");
s.push_str("# TYPE inklog_db_batch_size gauge\n");
s.push_str(&format!(
"inklog_db_batch_size{{sink=\"database\"}} {}\n",
self.db_batch_size.get()
));
s.push_str("# HELP inklog_db_batch_records_total Total records written by database batch flushes\n");
s.push_str("# TYPE inklog_db_batch_records_total counter\n");
s.push_str(&format!(
"inklog_db_batch_records_total{{sink=\"database\"}} {}\n",
self.db_batch_records_total.load(Ordering::Relaxed)
));
s.push_str("# HELP inklog_active_workers Current active worker threads\n");
s.push_str("# TYPE inklog_active_workers gauge\n");
s.push_str(&format!(
"inklog_active_workers {}\n",
self.active_workers.get()
));
let count = self.latency_count.load(Ordering::Relaxed);
let total = self.total_latency_us.load(Ordering::Relaxed);
let avg_latency = total.checked_div(count).unwrap_or(0);
s.push_str("# HELP inklog_avg_latency_us Average log processing latency in microseconds\n");
s.push_str("# TYPE inklog_avg_latency_us gauge\n");
s.push_str(&format!("inklog_avg_latency_us {}\n", avg_latency));
s.push_str("# HELP inklog_latency_p50_us P50 latency in microseconds\n");
s.push_str("# TYPE inklog_latency_p50_us gauge\n");
s.push_str(&format!(
"inklog_latency_p50_us {}\n",
self.latency_histogram.p50()
));
s.push_str("# HELP inklog_latency_p95_us P95 latency in microseconds\n");
s.push_str("# TYPE inklog_latency_p95_us gauge\n");
s.push_str(&format!(
"inklog_latency_p95_us {}\n",
self.latency_histogram.p95()
));
s.push_str("# HELP inklog_latency_p99_us P99 latency in microseconds\n");
s.push_str("# TYPE inklog_latency_p99_us gauge\n");
s.push_str(&format!(
"inklog_latency_p99_us {}\n",
self.latency_histogram.p99()
));
let uptime = self.uptime().as_secs();
if uptime > 0 {
s.push_str("# HELP inklog_uptime_seconds Uptime in seconds\n");
s.push_str("# TYPE inklog_uptime_seconds gauge\n");
s.push_str(&format!("inklog_uptime_seconds {}\n", uptime));
}
s.push_str("# HELP inklog_pool_hit_rate Pool hit rate percentage (0-100)\n");
s.push_str("# TYPE inklog_pool_hit_rate gauge\n");
s.push_str(&format!(
"inklog_pool_hit_rate {}\n",
self.pool_hit_rate.get()
));
s.push_str("# HELP inklog_sink_healthy Sink health status (1=healthy, 0=unhealthy)\n");
s.push_str("# TYPE inklog_sink_healthy gauge\n");
if let Ok(health_map) = self.sink_health.lock() {
for (name, health) in health_map.iter() {
let value = if health.status.is_operational() { 1 } else { 0 };
s.push_str(&format!(
"inklog_sink_healthy{{sink=\"{}\"}} {}\n",
name, value
));
}
}
s.push_str("# HELP inklog_latency_bucket Latency histogram bucket\n");
s.push_str("# TYPE inklog_latency_bucket counter\n");
let bounds = [1000, 5000, 10000, 50000, 100000, 500000, 1000000];
let buckets = self.latency_histogram.snapshot();
for (i, &bound) in bounds.iter().enumerate() {
if i < buckets.len() {
s.push_str(&format!(
"inklog_latency_bucket{{le=\"{}\"}} {}\n",
bound, buckets[i]
));
}
}
let total_count: u64 = buckets.iter().sum();
s.push_str(&format!(
"inklog_latency_bucket{{le=\"+Inf\"}} {}\n",
total_count
));
s
}
}
#[derive(Debug, Clone, PartialEq)]
pub enum FallbackState {
Active,
Fallback { target: String, reason: String },
Recovering { attempt: u32, delay_ms: u64 },
}
#[derive(Debug, Clone)]
pub struct FallbackConfig {
pub enabled: bool,
pub initial_delay_ms: u64,
pub max_delay_ms: u64,
pub max_retries: u32,
pub failure_threshold: u32,
}
impl Default for FallbackConfig {
fn default() -> Self {
Self {
enabled: true,
initial_delay_ms: 1000,
max_delay_ms: 60000,
max_retries: 10,
failure_threshold: 3,
}
}
}
#[derive(Debug)]
pub struct SinkHealthMonitor {
config: FallbackConfig,
fallback_states: Mutex<HashMap<String, FallbackState>>,
retry_counters: Mutex<HashMap<String, u32>>,
fallback_events: Mutex<Vec<FallbackEvent>>,
}
#[derive(Debug, Clone)]
pub struct FallbackEvent {
pub timestamp: chrono::DateTime<chrono::Utc>,
pub sink_name: String,
pub from_state: FallbackState,
pub to_state: FallbackState,
pub reason: String,
}
impl Default for SinkHealthMonitor {
fn default() -> Self {
Self::new(FallbackConfig::default())
}
}
impl SinkHealthMonitor {
pub fn new(config: FallbackConfig) -> Self {
Self {
config,
fallback_states: Mutex::new(HashMap::new()),
retry_counters: Mutex::new(HashMap::new()),
fallback_events: Mutex::new(Vec::new()),
}
}
pub fn with_defaults() -> Self {
Self::new(FallbackConfig::default())
}
pub fn check_and_fallback(
&self,
sink_name: &str,
is_healthy: bool,
error: Option<&str>,
) -> FallbackAction {
let mut states = self
.fallback_states
.lock()
.unwrap_or_else(|e| e.into_inner());
let mut retries = self
.retry_counters
.lock()
.unwrap_or_else(|e| e.into_inner());
let current_state = states
.get(sink_name)
.cloned()
.unwrap_or(FallbackState::Active);
if is_healthy {
self.handle_recovery(sink_name, ¤t_state, &mut states, &mut retries)
} else {
let error_msg = error.unwrap_or("Unknown error").to_string();
self.handle_failure(
sink_name,
&error_msg,
¤t_state,
&mut states,
&mut retries,
)
}
}
fn handle_recovery(
&self,
sink_name: &str,
current_state: &FallbackState,
states: &mut HashMap<String, FallbackState>,
retries: &mut HashMap<String, u32>,
) -> FallbackAction {
match current_state {
FallbackState::Active => {
FallbackAction::None
}
FallbackState::Fallback { target, reason: _ } => {
tracing::info!(
event = "sink_recovering",
sink = sink_name,
fallback_target = target,
"Sink {} 正在从 {} 恢复",
sink_name,
target
);
let attempt = retries.get(sink_name).cloned().unwrap_or(0) + 1;
retries.insert(sink_name.to_string(), attempt);
let delay_ms = self
.config
.initial_delay_ms
.saturating_mul(2_u64.pow(attempt.min(10)))
.min(self.config.max_delay_ms);
states.insert(
sink_name.to_string(),
FallbackState::Recovering { attempt, delay_ms },
);
self.log_event(
sink_name,
current_state.clone(),
states
.get(sink_name)
.cloned()
.unwrap_or(current_state.clone()),
format!("尝试恢复,延迟 {}ms", delay_ms),
);
FallbackAction::AttemptRecovery {
sink_name: sink_name.to_string(),
attempt,
delay_ms,
}
}
FallbackState::Recovering {
attempt: _,
delay_ms,
} => {
FallbackAction::Wait {
sink_name: sink_name.to_string(),
remaining_ms: *delay_ms,
}
}
}
}
fn handle_failure(
&self,
sink_name: &str,
error: &str,
current_state: &FallbackState,
states: &mut HashMap<String, FallbackState>,
retries: &mut HashMap<String, u32>,
) -> FallbackAction {
let failure_count = retries.get(sink_name).cloned().unwrap_or(0) + 1;
retries.insert(sink_name.to_string(), failure_count);
let (fallback_target, action) = self.determine_fallback_target(sink_name, error);
if !self.config.enabled {
tracing::warn!(
event = "sink_failed_disabled_fallback",
sink = sink_name,
error = error,
"Sink {} 故障但自动降级已禁用",
sink_name
);
return FallbackAction::None;
}
if failure_count >= self.config.failure_threshold {
let new_state = FallbackState::Fallback {
target: fallback_target.clone(),
reason: error.to_string(),
};
tracing::warn!(
event = "sink_fallback_triggered",
sink = sink_name,
fallback_target = fallback_target,
error = error,
failure_count = failure_count,
"Sink {} 降级到 {},原因: {}",
sink_name,
fallback_target,
error
);
states.insert(sink_name.to_string(), new_state.clone());
self.log_event(
sink_name,
current_state.clone(),
new_state,
error.to_string(),
);
retries.insert(sink_name.to_string(), 0);
action
} else {
tracing::warn!(
event = "sink_failure_warning",
sink = sink_name,
error = error,
failure_count = failure_count,
threshold = self.config.failure_threshold,
"Sink {} 连续第 {} 次故障",
sink_name,
failure_count
);
FallbackAction::Retry {
sink_name: sink_name.to_string(),
attempt: failure_count,
error: error.to_string(),
}
}
}
fn determine_fallback_target(&self, sink_name: &str, error: &str) -> (String, FallbackAction) {
match sink_name {
"database" => {
(
"file".to_string(),
FallbackAction::Fallback {
sink_name: sink_name.to_string(),
target: "file".to_string(),
reason: format!("Database 故障: {}", error),
},
)
}
"file" => {
let lower_error = error.to_lowercase();
if lower_error.contains("disk")
|| lower_error.contains("space")
|| lower_error.contains("full")
{
(
"console".to_string(),
FallbackAction::Fallback {
sink_name: sink_name.to_string(),
target: "console".to_string(),
reason: format!("磁盘空间不足: {}", error),
},
)
} else {
(
"console".to_string(),
FallbackAction::Fallback {
sink_name: sink_name.to_string(),
target: "console".to_string(),
reason: format!("FileSink 故障: {}", error),
},
)
}
}
_ => {
(
"console".to_string(),
FallbackAction::Fallback {
sink_name: sink_name.to_string(),
target: "console".to_string(),
reason: format!("未知故障: {}", error),
},
)
}
}
}
pub fn handle_encryption_error(&self, sink_name: &str, error: &str) -> FallbackAction {
tracing::warn!(
event = "encryption_error_fallback",
sink = sink_name,
error = error,
"加密密钥错误,降级为明文写入"
);
let mut states = self
.fallback_states
.lock()
.unwrap_or_else(|e| e.into_inner());
let current_state = states
.get(sink_name)
.cloned()
.unwrap_or(FallbackState::Active);
let new_state = FallbackState::Fallback {
target: "plaintext".to_string(),
reason: format!("加密密钥错误: {}", error),
};
states.insert(sink_name.to_string(), new_state.clone());
self.log_event(sink_name, current_state, new_state, error.to_string());
FallbackAction::Fallback {
sink_name: sink_name.to_string(),
target: "plaintext".to_string(),
reason: format!("明文写入(加密错误): {}", error),
}
}
pub fn confirm_recovery(&self, sink_name: &str) {
let mut states = self
.fallback_states
.lock()
.unwrap_or_else(|e| e.into_inner());
let mut retries = self
.retry_counters
.lock()
.unwrap_or_else(|e| e.into_inner());
let current_state = states
.get(sink_name)
.cloned()
.unwrap_or(FallbackState::Active);
if !matches!(current_state, FallbackState::Active) {
tracing::info!(
event = "sink_recovery_confirmed",
sink = sink_name,
"Sink {} 恢复成功,已切回正常模式",
sink_name
);
states.insert(sink_name.to_string(), FallbackState::Active);
retries.remove(sink_name);
self.log_event(
sink_name,
current_state,
FallbackState::Active,
"恢复成功".to_string(),
);
}
}
pub fn get_fallback_state(&self, sink_name: &str) -> FallbackState {
self.fallback_states
.lock()
.unwrap_or_else(|e| e.into_inner())
.get(sink_name)
.cloned()
.unwrap_or(FallbackState::Active)
}
pub fn is_any_in_fallback(&self) -> bool {
self.fallback_states
.lock()
.unwrap_or_else(|e| e.into_inner())
.values()
.any(|state| matches!(state, FallbackState::Fallback { .. }))
}
pub fn get_fallback_events(&self, limit: usize) -> Vec<FallbackEvent> {
let events = self
.fallback_events
.lock()
.unwrap_or_else(|e| e.into_inner());
events.iter().rev().take(limit).cloned().collect()
}
pub fn get_fallback_stats(&self) -> FallbackStats {
let states = self
.fallback_states
.lock()
.unwrap_or_else(|e| e.into_inner());
let events = self
.fallback_events
.lock()
.unwrap_or_else(|e| e.into_inner());
let active_fallbacks = states
.values()
.filter(|s| matches!(s, FallbackState::Fallback { .. }))
.count();
let recovering = states
.values()
.filter(|s| matches!(s, FallbackState::Recovering { .. }))
.count();
let fallback_events_count = events.len();
FallbackStats {
active_fallbacks,
recovering,
total_fallback_events: fallback_events_count,
}
}
fn log_event(&self, sink_name: &str, from: FallbackState, to: FallbackState, reason: String) {
let event = FallbackEvent {
timestamp: chrono::Utc::now(),
sink_name: sink_name.to_string(),
from_state: from,
to_state: to,
reason,
};
let mut events = self
.fallback_events
.lock()
.unwrap_or_else(|e| e.into_inner());
events.push(event);
if events.len() > 100 {
events.remove(0);
}
}
pub fn reset(&self) {
let mut states = self
.fallback_states
.lock()
.unwrap_or_else(|e| e.into_inner());
let mut retries = self
.retry_counters
.lock()
.unwrap_or_else(|e| e.into_inner());
let mut events = self
.fallback_events
.lock()
.unwrap_or_else(|e| e.into_inner());
states.clear();
retries.clear();
events.clear();
}
}
#[derive(Debug, Clone)]
pub enum FallbackAction {
None,
Retry {
sink_name: String,
attempt: u32,
error: String,
},
Fallback {
sink_name: String,
target: String,
reason: String,
},
AttemptRecovery {
sink_name: String,
attempt: u32,
delay_ms: u64,
},
Wait {
sink_name: String,
remaining_ms: u64,
},
LocalQueue { sink_name: String, reason: String },
}
impl FallbackAction {
pub fn requires_action(&self) -> bool {
!matches!(self, FallbackAction::None)
}
pub fn sink_name(&self) -> Option<&str> {
match self {
FallbackAction::None => None,
FallbackAction::Retry { sink_name, .. } => Some(sink_name),
FallbackAction::Fallback { sink_name, .. } => Some(sink_name),
FallbackAction::AttemptRecovery { sink_name, .. } => Some(sink_name),
FallbackAction::Wait { sink_name, .. } => Some(sink_name),
FallbackAction::LocalQueue { sink_name, .. } => Some(sink_name),
}
}
}
#[derive(Debug, Clone, Default)]
pub struct FallbackStats {
pub active_fallbacks: usize,
pub recovering: usize,
pub total_fallback_events: usize,
}
#[cfg(test)]
mod sink_health_monitor_tests {
use super::*;
#[test]
fn test_fallback_state_operations() {
let monitor = SinkHealthMonitor::with_defaults();
assert_eq!(
monitor.get_fallback_state("database"),
FallbackState::Active
);
assert!(!monitor.is_any_in_fallback());
let action = monitor.check_and_fallback("database", false, Some("Connection refused"));
assert!(action.requires_action());
assert!(matches!(action, FallbackAction::Retry { .. }));
for _ in 0..3 {
let _ = monitor.check_and_fallback("database", false, Some("Connection refused"));
}
let state = monitor.get_fallback_state("database");
assert!(matches!(state, FallbackState::Fallback { target, .. } if target == "file"));
let action = monitor.check_and_fallback("database", true, None);
assert!(matches!(action, FallbackAction::AttemptRecovery { .. }));
monitor.confirm_recovery("database");
assert_eq!(
monitor.get_fallback_state("database"),
FallbackState::Active
);
}
#[test]
fn test_file_sink_disk_full_fallback() {
let monitor = SinkHealthMonitor::with_defaults();
for _ in 0..3 {
let _ = monitor.check_and_fallback("file", false, Some("Disk is full"));
}
let state = monitor.get_fallback_state("file");
assert!(matches!(
state,
FallbackState::Fallback { target, .. } if target == "console"
));
}
#[test]
fn test_encryption_error_fallback() {
let monitor = SinkHealthMonitor::with_defaults();
let action = monitor.handle_encryption_error("file", "Invalid key");
assert!(matches!(action, FallbackAction::Fallback { target, .. } if target == "plaintext"));
}
#[test]
fn test_fallback_stats() {
let monitor = SinkHealthMonitor::with_defaults();
let stats = monitor.get_fallback_stats();
assert_eq!(stats.active_fallbacks, 0);
assert_eq!(stats.recovering, 0);
}
#[test]
fn test_fallback_events() {
let monitor = SinkHealthMonitor::with_defaults();
for _ in 0..3 {
let _ = monitor.check_and_fallback("database", false, Some("Error"));
}
let events = monitor.get_fallback_events(10);
assert!(!events.is_empty());
assert_eq!(events[0].sink_name, "database");
}
#[test]
fn test_disabled_fallback() {
let config = FallbackConfig {
enabled: false,
..Default::default()
};
let monitor = SinkHealthMonitor::new(config);
let action = monitor.check_and_fallback("database", false, Some("Error"));
assert!(!action.requires_action());
}
#[test]
fn test_reset() {
let monitor = SinkHealthMonitor::with_defaults();
for _ in 0..3 {
let _ = monitor.check_and_fallback("database", false, Some("Error"));
}
assert!(monitor.is_any_in_fallback());
monitor.reset();
assert!(!monitor.is_any_in_fallback());
assert_eq!(
monitor.get_fallback_state("database"),
FallbackState::Active
);
}
#[test]
fn test_default_sink_health_monitor() {
let monitor = SinkHealthMonitor::default();
assert_eq!(monitor.get_fallback_state("any"), FallbackState::Active);
}
#[test]
fn test_recovery_from_active_state() {
let monitor = SinkHealthMonitor::with_defaults();
let action = monitor.check_and_fallback("database", true, None);
assert!(matches!(action, FallbackAction::None));
}
#[test]
fn test_recovery_wait_from_recovering_state() {
let monitor = SinkHealthMonitor::with_defaults();
for _ in 0..3 {
let _ = monitor.check_and_fallback("database", false, Some("Error"));
}
let action = monitor.check_and_fallback("database", true, None);
assert!(matches!(action, FallbackAction::AttemptRecovery { .. }));
let action = monitor.check_and_fallback("database", true, None);
assert!(matches!(action, FallbackAction::Wait { .. }));
}
#[test]
fn test_determine_fallback_target_file_space() {
let monitor = SinkHealthMonitor::with_defaults();
for _ in 0..3 {
let _ = monitor.check_and_fallback("file", false, Some("No space left"));
}
let state = monitor.get_fallback_state("file");
assert!(matches!(state, FallbackState::Fallback { target, .. } if target == "console"));
}
#[test]
fn test_determine_fallback_target_file_full() {
let monitor = SinkHealthMonitor::with_defaults();
for _ in 0..3 {
let _ = monitor.check_and_fallback("file", false, Some("storage full"));
}
let state = monitor.get_fallback_state("file");
assert!(matches!(state, FallbackState::Fallback { target, .. } if target == "console"));
}
#[test]
fn test_determine_fallback_target_file_other_error() {
let monitor = SinkHealthMonitor::with_defaults();
for _ in 0..3 {
let _ = monitor.check_and_fallback("file", false, Some("permission denied"));
}
let state = monitor.get_fallback_state("file");
assert!(matches!(state, FallbackState::Fallback { target, .. } if target == "console"));
}
#[test]
fn test_determine_fallback_target_unknown_sink() {
let monitor = SinkHealthMonitor::with_defaults();
for _ in 0..3 {
let _ = monitor.check_and_fallback("custom_sink", false, Some("Unknown failure"));
}
let state = monitor.get_fallback_state("custom_sink");
assert!(matches!(state, FallbackState::Fallback { target, .. } if target == "console"));
}
#[test]
fn test_fallback_action_sink_name_all_variants() {
let fallback = FallbackAction::Fallback {
sink_name: "db".to_string(),
target: "file".to_string(),
reason: "err".to_string(),
};
assert_eq!(fallback.sink_name(), Some("db"));
let recovery = FallbackAction::AttemptRecovery {
sink_name: "db".to_string(),
attempt: 1,
delay_ms: 100,
};
assert_eq!(recovery.sink_name(), Some("db"));
let wait = FallbackAction::Wait {
sink_name: "db".to_string(),
remaining_ms: 100,
};
assert_eq!(wait.sink_name(), Some("db"));
let queue = FallbackAction::LocalQueue {
sink_name: "s3".to_string(),
reason: "err".to_string(),
};
assert_eq!(queue.sink_name(), Some("s3"));
}
#[test]
fn test_log_event_truncation() {
let monitor = SinkHealthMonitor::with_defaults();
for _ in 0..(102 * 3) {
let _ = monitor.check_and_fallback("database", false, Some("Error"));
}
let events = monitor.get_fallback_events(200);
assert_eq!(events.len(), 100);
}
}
#[cfg(test)]
mod metrics_tests {
use super::*;
#[test]
fn test_gauge_new() {
let gauge = Gauge::new(100);
assert_eq!(gauge.get(), 100);
}
#[test]
fn test_gauge_set() {
let gauge = Gauge::new(0);
gauge.set(42);
assert_eq!(gauge.get(), 42);
}
#[test]
fn test_gauge_inc() {
let gauge = Gauge::new(0);
gauge.inc();
gauge.inc();
assert_eq!(gauge.get(), 2);
}
#[test]
fn test_gauge_dec() {
let gauge = Gauge::new(10);
gauge.dec();
gauge.dec();
assert_eq!(gauge.get(), 8);
}
#[test]
fn test_histogram_new() {
let histogram = Histogram::new(vec![100, 500, 1000]);
assert_eq!(histogram.buckets.len(), 4);
}
#[test]
fn test_histogram_record() {
let histogram = Histogram::new(vec![100, 500, 1000]);
histogram.record(50);
histogram.record(200);
histogram.record(600);
histogram.record(1500);
let snapshot = histogram.snapshot();
assert_eq!(snapshot[0], 1);
assert_eq!(snapshot[1], 1);
assert_eq!(snapshot[2], 1);
assert_eq!(snapshot[3], 1);
}
#[test]
fn test_histogram_percentile_empty_bounds_edge() {
let histogram = Histogram::new(vec![100, 500, 1000]);
histogram.record(50);
let p50 = histogram.percentile(50.0);
assert_eq!(p50, 100);
}
#[test]
fn test_histogram_percentile_returns_last_bound_when_target_exceeds_total() {
let hist = Histogram::new(vec![100, 500, 1000]);
hist.record(50);
let result = hist.percentile(200.0);
assert_eq!(result, 1000);
}
#[test]
fn test_gauge_f64_new_and_set() {
let gauge = GaugeF64::new(1.5);
assert_eq!(gauge.get(), 1.5);
gauge.set(3.15);
assert_eq!(gauge.get(), 3.15);
}
#[test]
fn test_gauge_f64_mutex_poisoning_recovery_set() {
let gauge = std::sync::Arc::new(GaugeF64::new(0.0));
let gauge_clone = gauge.clone();
let handle = std::thread::spawn(move || {
let _guard = gauge_clone.value.lock().unwrap();
panic!("Intentional panic to poison mutex");
});
let _ = handle.join();
gauge.set(42.0);
assert_eq!(gauge.get(), 42.0);
}
#[test]
fn test_gauge_f64_mutex_poisoning_recovery_get() {
let gauge = std::sync::Arc::new(GaugeF64::new(10.0));
let gauge_clone = gauge.clone();
let handle = std::thread::spawn(move || {
let _guard = gauge_clone.value.lock().unwrap();
panic!("Intentional panic to poison mutex");
});
let _ = handle.join();
let val = gauge.get();
assert_eq!(val, 10.0);
}
#[test]
fn test_sink_health_healthy() {
let health = SinkHealth::healthy();
assert!(matches!(health.status, SinkStatus::Healthy));
assert!(health.last_error.is_none());
assert_eq!(health.consecutive_failures, 0);
}
#[test]
fn test_sink_health_unhealthy() {
let health = SinkHealth::unhealthy("Connection refused".to_string());
assert!(matches!(health.status, SinkStatus::Unhealthy { .. }));
assert!(health.last_error.is_some());
assert_eq!(health.consecutive_failures, 1);
}
#[test]
fn test_metrics_new() {
let metrics = Metrics::new();
assert_eq!(metrics.logs_written(), 0);
assert_eq!(metrics.logs_dropped(), 0);
assert_eq!(metrics.sink_errors(), 0);
}
#[test]
fn test_metrics_record_log_written() {
let metrics = Metrics::new();
metrics.inc_logs_written();
assert_eq!(metrics.logs_written(), 1);
metrics.inc_logs_written();
assert_eq!(metrics.logs_written(), 2);
}
#[test]
fn test_metrics_record_log_dropped() {
let metrics = Metrics::new();
metrics.inc_logs_dropped();
assert_eq!(metrics.logs_dropped(), 1);
}
#[test]
fn test_metrics_record_sink_error() {
let metrics = Metrics::new();
metrics.inc_sink_error();
assert_eq!(metrics.sink_errors(), 1);
}
#[test]
fn test_metrics_record_db_batch() {
let metrics = Metrics::new();
metrics.set_db_batch_size(8);
metrics.add_db_batch_records_total(8);
assert_eq!(metrics.db_batch_size(), 8);
assert_eq!(metrics.db_batch_records_total(), 8);
}
#[test]
fn test_metrics_record_channel_blocked() {
let metrics = Metrics::new();
metrics.inc_channel_blocked();
assert_eq!(metrics.channel_blocked(), 1);
}
#[test]
fn test_metrics_record_latency() {
let metrics = Metrics::new();
metrics.record_latency(Duration::from_micros(100));
metrics.record_latency(Duration::from_micros(200));
metrics.record_latency(Duration::from_micros(300));
let latency = metrics.total_latency_us.load(Ordering::Relaxed);
assert!(latency >= 100 + 200 + 300);
}
#[test]
fn test_metrics_update_sink_health() {
let metrics = Metrics::new();
metrics.update_sink_health("file", true, None);
let health = metrics.sink_health();
assert!(health.contains_key("file"));
}
#[test]
fn test_metrics_active_workers() {
let metrics = Metrics::new();
metrics.active_workers.set(4);
assert_eq!(metrics.active_workers(), 4);
}
#[test]
fn test_metrics_logs_written() {
let metrics = Metrics::new();
metrics.inc_logs_written();
metrics.inc_logs_written();
assert_eq!(metrics.logs_written(), 2);
}
#[test]
fn test_metrics_logs_dropped() {
let metrics = Metrics::new();
metrics.inc_logs_dropped();
assert_eq!(metrics.logs_dropped(), 1);
}
#[test]
fn test_metrics_export_prometheus_db_batch() {
let metrics = Metrics::new();
metrics.set_db_batch_size(5);
metrics.add_db_batch_records_total(12);
let output = metrics.export_prometheus();
assert!(output.contains("inklog_db_batch_size{sink=\"database\"} 5"));
assert!(output.contains("inklog_db_batch_records_total{sink=\"database\"} 12"));
}
#[test]
fn test_metrics_channel_blocked() {
let metrics = Metrics::new();
metrics.inc_channel_blocked();
assert_eq!(metrics.channel_blocked(), 1);
}
#[test]
fn test_metrics_sink_errors() {
let metrics = Metrics::new();
metrics.inc_sink_error();
assert_eq!(metrics.sink_errors(), 1);
}
#[test]
fn test_metrics_uptime() {
let metrics = Metrics::new();
let uptime = metrics.uptime();
assert!(uptime <= std::time::Duration::from_secs(60));
}
#[test]
fn test_gauge_f64_set_and_get() {
let gauge = GaugeF64::new(0.5);
assert!((gauge.get() - 0.5).abs() < f64::EPSILON);
gauge.set(99.9);
assert!((gauge.get() - 99.9).abs() < f64::EPSILON);
}
#[test]
fn test_gauge_f64_pool_hit_rate() {
let metrics = Metrics::new();
metrics.set_pool_hit_rate(85.5);
assert!((metrics.pool_hit_rate() - 85.5).abs() < f64::EPSILON);
metrics.set_pool_hit_rate(0.0);
assert!((metrics.pool_hit_rate() - 0.0).abs() < f64::EPSILON);
}
#[test]
fn test_sink_status_is_operational() {
assert!(SinkStatus::Healthy.is_operational());
assert!(
SinkStatus::Degraded {
reason: "slow".to_string()
}
.is_operational()
);
assert!(
!SinkStatus::Unhealthy {
error: "crashed".to_string()
}
.is_operational()
);
assert!(!SinkStatus::NotStarted.is_operational());
}
#[test]
fn test_fallback_action_sink_name() {
let none = FallbackAction::None;
assert_eq!(none.sink_name(), None);
assert!(!none.requires_action());
let retry = FallbackAction::Retry {
sink_name: "db".to_string(),
attempt: 1,
error: "err".to_string(),
};
assert_eq!(retry.sink_name(), Some("db"));
assert!(retry.requires_action());
}
#[test]
fn test_histogram_percentile_empty() {
let histogram = Histogram::new(vec![100, 500, 1000]);
assert_eq!(histogram.percentile(50.0), 0);
assert_eq!(histogram.p50(), 0);
assert_eq!(histogram.p95(), 0);
assert_eq!(histogram.p99(), 0);
}
#[test]
fn test_histogram_percentile() {
let histogram = Histogram::new(vec![100, 500, 1000]);
histogram.record(50);
histogram.record(50);
histogram.record(50);
histogram.record(200);
histogram.record(200);
histogram.record(600);
histogram.record(2000);
assert_eq!(histogram.p50(), 500);
}
#[test]
fn test_metrics_export_prometheus_format() {
let metrics = Metrics::new();
metrics.inc_logs_written();
metrics.inc_logs_written();
metrics.inc_sink_error();
metrics.set_db_batch_size(10);
let output = metrics.export_prometheus();
assert!(output.contains("# TYPE inklog_logs_written_total counter"));
assert!(output.contains("inklog_logs_written_total 2"));
assert!(output.contains("# TYPE inklog_sink_errors_total counter"));
assert!(output.contains("inklog_sink_errors_total 1"));
assert!(output.contains("inklog_db_batch_size{sink=\"database\"} 10"));
}
#[test]
fn test_fallback_config_default() {
let config = FallbackConfig::default();
assert!(config.enabled);
assert_eq!(config.initial_delay_ms, 1000);
assert_eq!(config.max_delay_ms, 60000);
assert_eq!(config.max_retries, 10);
assert_eq!(config.failure_threshold, 3);
}
#[test]
fn test_fallback_event_structure() {
let event = FallbackEvent {
timestamp: chrono::Utc::now(),
sink_name: "database".to_string(),
from_state: FallbackState::Active,
to_state: FallbackState::Fallback {
target: "file".to_string(),
reason: "disk error".to_string(),
},
reason: "Disk is full".to_string(),
};
assert_eq!(event.sink_name, "database");
assert!(matches!(event.from_state, FallbackState::Active));
}
#[test]
fn test_sink_health_default() {
let health = SinkHealth::default();
assert!(matches!(health.status, SinkStatus::NotStarted));
assert!(health.last_error.is_none());
assert_eq!(health.consecutive_failures, 0);
}
#[test]
fn test_lock_contention_metrics() {
let metrics = Metrics::new();
assert_eq!(metrics.lock_contention(), 0);
metrics.inc_lock_contention();
metrics.inc_lock_contention();
assert_eq!(metrics.lock_contention(), 2);
}
#[test]
fn test_update_sink_health_unhealthy_no_error() {
let metrics = Metrics::new();
metrics.update_sink_health("file", false, None);
let health = metrics.sink_health();
let h = health.get("file").expect("sink health should exist");
assert!(matches!(&h.status, SinkStatus::Unhealthy { error } if error == "Unknown error"));
assert_eq!(h.consecutive_failures, 1);
}
#[test]
fn test_update_sink_health_consecutive_failures() {
let metrics = Metrics::new();
metrics.update_sink_health("file", false, Some("err1".to_string()));
let health = metrics.sink_health();
let h = health.get("file").expect("should exist");
assert_eq!(h.consecutive_failures, 1);
metrics.update_sink_health("file", false, Some("err2".to_string()));
let health = metrics.sink_health();
let h = health.get("file").expect("should exist");
assert_eq!(h.consecutive_failures, 2);
metrics.update_sink_health("file", true, None);
let health = metrics.sink_health();
let h = health.get("file").expect("should exist");
assert_eq!(h.consecutive_failures, 0);
}
#[test]
fn test_sink_started() {
let metrics = Metrics::new();
metrics.sink_started("console");
let health = metrics.sink_health();
let h = health.get("console").expect("sink should be started");
assert!(matches!(h.status, SinkStatus::Healthy));
assert_eq!(h.consecutive_failures, 0);
assert!(h.last_error.is_none());
}
#[test]
fn test_sink_degraded() {
let metrics = Metrics::new();
metrics.sink_degraded("file", "slow disk".to_string());
let health = metrics.sink_health();
let h = health.get("file").expect("sink should exist");
assert!(matches!(&h.status, SinkStatus::Degraded { reason } if reason == "slow disk"));
assert_eq!(h.last_error, Some("slow disk".to_string()));
}
#[test]
fn test_get_status_all_healthy() {
let metrics = Metrics::new();
metrics.sink_started("console");
metrics.sink_started("file");
let status = metrics.get_status(0, 100);
assert!(matches!(status.overall_status, SinkStatus::Healthy));
}
#[test]
fn test_get_status_with_unhealthy() {
let metrics = Metrics::new();
metrics.sink_started("console");
metrics.update_sink_health("file", false, Some("disk error".to_string()));
let status = metrics.get_status(0, 100);
assert!(
matches!(status.overall_status, SinkStatus::Unhealthy { error } if error.contains("disk error"))
);
}
#[test]
fn test_get_status_with_degraded() {
let metrics = Metrics::new();
metrics.sink_started("console");
metrics.sink_degraded("file", "slow disk".to_string());
let status = metrics.get_status(0, 100);
assert!(
matches!(status.overall_status, SinkStatus::Degraded { reason } if reason.contains("slow disk"))
);
}
#[test]
fn test_export_prometheus_with_uptime() {
let metrics = Metrics::new();
std::thread::sleep(std::time::Duration::from_secs(1));
let output = metrics.export_prometheus();
assert!(output.contains("inklog_uptime_seconds"));
}
#[test]
fn test_export_prometheus_with_sink_health() {
let metrics = Metrics::new();
metrics.sink_started("console");
metrics.update_sink_health("file", false, Some("err".to_string()));
let output = metrics.export_prometheus();
assert!(output.contains("inklog_sink_healthy{sink=\"console\"} 1"));
assert!(output.contains("inklog_sink_healthy{sink=\"file\"} 0"));
}
#[test]
fn test_histogram_percentile_last_bucket() {
let histogram = Histogram::new(vec![100, 500, 1000]);
histogram.record(2000);
assert_eq!(histogram.p99(), 1000);
}
#[test]
fn test_get_status_with_not_started_sink_returns_healthy() {
let metrics = Metrics::new();
if let Ok(mut map) = metrics.sink_health.lock() {
map.insert("not_started_sink".to_string(), SinkHealth::default());
}
let status = metrics.get_status(0, 100);
assert!(
matches!(status.overall_status, SinkStatus::Healthy),
"expected Healthy, got {:?}",
status.overall_status
);
}
}