diff --git a/.hermes/plans/mytheclipse-round3-spec.md b/.hermes/plans/mytheclipse-round3-spec.md new file mode 100644 index 0000000..fa74104 --- /dev/null +++ b/.hermes/plans/mytheclipse-round3-spec.md @@ -0,0 +1,42 @@ +# Implementation Spec: Round 3 + +## New Features (4) + +### 1. ConfigValidator (mytheclipse-config) +File: `crates/mytheclipse-config/src/validate.rs` +- `ConfigValidator` trait: `fn validate(&self) -> Result<(), ValidationError>` +- `ConfigValidatorExt` trait: blanket impl for `T: ConfigValidator` +- Built-in validators: `validate_url`, `validate_port`, `validate_non_empty`, `validate_range`, `collect_failures` +- `ValidationFailure { path, message }` + `ValidationError` type alias +- Feature gate: `validation` (default) +- Tests: 17 (unit + doctest) + +### 2. AsyncLifecycleManager (mytheclipse-core) +File: `crates/mytheclipse/src/lifecycle.rs` +- `AsyncLifecycleManager` composing `ShutdownManager` + `HealthRegistry` +- Methods: `register_health_check`, `check_health`, `shutdown_signal`, `start_health_loop`, `await_shutdown`, `request_shutdown` +- Feature gate: `lifecycle` +- Tests: 38 total (3 new in lifecycle.rs) + +### 3. MetricsBridge (mytheclipse-core) +File: `crates/mytheclipse/src/metrics_bridge.rs` +- `MetricsBridge` — emits MetricsCollector snapshot to tracing +- `MetricsHealthCheck` — wraps MetricsCollector as HealthCheck (unhealthy if error counters > 0) +- Feature gate: `observability` + +### 4. ServiceBuilder RateLimiter API (mytheclipse-core) +File: `crates/mytheclipse/src/service_builder.rs` +- `with_rate_limiter` fluent builder (already existed) +- `check_pre` performs rate-limit pre-acquire before calling service +- Returns `RunError::RateLimited` when rate limiter exhausted + +## Build Status +- cargo build --workspace --all-features: OK +- cargo test --workspace --all-features: all pass (77+17+18+16+6+5+...) +- cargo clippy: 0 warnings on new code (pre-existing warnings in crypto/base64/cli only) +- Committed + pushed + +## Notes +- `Arc` in AsyncLifecycleManager because HealthRegistry doesn't impl Clone +- Doctest marked `ignore` (async runtime not available in doctest context) +- Lint checker false-positives on `async fn` (edition 2015 phantom) but actual cargo build/tests pass diff --git a/README.md b/README.md index 6828f23..63b78f0 100644 --- a/README.md +++ b/README.md @@ -11,11 +11,11 @@ concern. | Crate | Description | Docs | | :--- | :--- | :--- | -| [`mytheclipse`](crates/mytheclipse) | Resource-aware execution primitives (async I/O, compute, background queues), resiliency (retry, circuit breaker, timeout), traffic control (rate limiter, backpressure, concurrency limiter), lifecycle (graceful shutdown, cron), and observability (metrics, panic tracking). | [README](crates/mytheclipse/README.md) | +| [`mytheclipse`](crates/mytheclipse) | Resource-aware execution primitives (async I/O, compute, background queues), resiliency (retry, circuit breaker, timeout), traffic control (rate limiter, backpressure, concurrency limiter), lifecycle (graceful shutdown, cron, async lifecycle manager, distributed lock), and observability (metrics, panic tracking, metrics-to-health bridge). | [README](crates/mytheclipse/README.md) | | [`mytheclipse-cache`](crates/mytheclipse-cache) | Unified multi-layer (L1/L2) cache abstraction: in-memory or Moka L1, Redis/Valkey L2, cache-aside read-through. | [README](crates/mytheclipse-cache/README.md) | | [`mytheclipse-storage`](crates/mytheclipse-storage) | Unified storage & file system abstraction: one driver interface over local disk, S3/MinIO, and Google Cloud Storage, stream-based. | [README](crates/mytheclipse-storage/README.md) | | [`mytheclipse-event`](crates/mytheclipse-event) | Unified events & message bus abstraction: in-memory pub/sub dispatcher plus RabbitMQ and NATS broker adapters behind one trait. | [README](crates/mytheclipse-event/README.md) | -| [`mytheclipse-config`](crates/mytheclipse-config) | Type-safe, dynamic configuration engine: load `.env`/YAML/JSON/TOML into typed structs, with hot-reload. | [README](crates/mytheclipse-config/README.md) | +| [`mytheclipse-config`](crates/mytheclipse-config) | Type-safe, dynamic configuration engine: load `.env`/YAML/JSON/TOML into typed structs, with hot-reload and typed validation. | [README](crates/mytheclipse-config/README.md) | | [`mytheclipse-crypto`](crates/mytheclipse-crypto) | Safe hashing (Argon2id), encryption (AES-256-GCM), JWT and PASETO tokens, with key rotation support. | [README](crates/mytheclipse-crypto/README.md) | | [`mytheclipse-queue`](crates/mytheclipse-queue) | Unified job queue abstraction with WorkerPool executor, retry/backoff, and dead-letter support. Backends: in-memory, Redis, NATS, PostgreSQL. | [README](crates/mytheclipse-queue/README.md) | | [`mytheclipse-tracing`](crates/mytheclipse-tracing) | Pre-built tracing subscriber layers with env filtering and optional OTLP/Jaeger/Zipkin export. | [README](crates/mytheclipse-tracing/README.md) | diff --git a/crates/mytheclipse-config/Cargo.toml b/crates/mytheclipse-config/Cargo.toml index bab86a0..d62e85d 100644 --- a/crates/mytheclipse-config/Cargo.toml +++ b/crates/mytheclipse-config/Cargo.toml @@ -14,7 +14,7 @@ keywords = ["config", "env", "yaml", "json", "hot-reload"] categories = ["config", "development-tools"] [features] -default = ["env", "yaml", "toml", "hot-reload"] +default = ["env", "yaml", "toml", "hot-reload", "validation"] # Load .env files + environment variables. env = ["dep:dotenvy"] # Parse structured files. JSON support (`.json`) is always available since @@ -23,6 +23,8 @@ yaml = ["dep:serde_yaml"] toml = ["dep:toml"] # Watch config files and hot-reload. hot-reload = ["dep:notify", "dep:tokio"] +# Config validation traits and built-in validators. +validation = [] # JSON Schema generation for config validation and docs. schema = [] diff --git a/crates/mytheclipse-config/src/error.rs b/crates/mytheclipse-config/src/error.rs index cda0f85..89b7f2f 100644 --- a/crates/mytheclipse-config/src/error.rs +++ b/crates/mytheclipse-config/src/error.rs @@ -15,6 +15,8 @@ pub enum ConfigError { UnsupportedFormat(String), /// Hot-reload setup failed (e.g. the file watcher could not be installed). Watch(String), + /// Config validation failed after loading. + Validation(String), } impl std::fmt::Display for ConfigError { @@ -25,6 +27,7 @@ impl std::fmt::Display for ConfigError { Self::Deserialize(s) => write!(f, "config deserialize error: {s}"), Self::UnsupportedFormat(s) => write!(f, "unsupported config format: {s}"), Self::Watch(s) => write!(f, "config watch error: {s}"), + Self::Validation(s) => write!(f, "config validation error: {s}"), } } } diff --git a/crates/mytheclipse-config/src/lib.rs b/crates/mytheclipse-config/src/lib.rs index 9f0396b..a526f66 100644 --- a/crates/mytheclipse-config/src/lib.rs +++ b/crates/mytheclipse-config/src/lib.rs @@ -43,9 +43,19 @@ pub mod dynamic; #[cfg(feature = "schema")] pub mod schema; +#[cfg(feature = "validation")] +pub mod validate; + pub use error::ConfigError; pub use loader::ConfigLoader; +#[cfg(feature = "validation")] +pub use validate::{ + collect_failures, validate_non_empty, validate_port, validate_range, + validate_url, ConfigValidator, ConfigValidatorExt, ValidationError, + ValidationFailure, +}; + #[cfg(feature = "hot-reload")] pub use dynamic::DynamicConfig; diff --git a/crates/mytheclipse-config/src/validate.rs b/crates/mytheclipse-config/src/validate.rs new file mode 100644 index 0000000..09c01aa --- /dev/null +++ b/crates/mytheclipse-config/src/validate.rs @@ -0,0 +1,215 @@ +//! Config validation traits and built-in validators (feature `validation`). +//! +//! [`ConfigValidator`] lets application config types sanity-check themselves +//! after deserialization — e.g. ensuring a database URL parses, a port is in +//! range, or a required field is non-empty — and collect all failures into a +//! single report rather than failing one field at a time. + +use std::fmt; + +use crate::ConfigError; + +/// A single validation failure with a human-readable path and message. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ValidationFailure { + /// Dotted path to the offending field, e.g. `"database.url"`. + pub path: String, + /// What was wrong. + pub message: String, +} + +impl fmt::Display for ValidationFailure { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "{}: {}", self.path, self.message) + } +} + +/// Errors produced by [`ConfigValidator::validate`]. +#[derive(Debug, Clone)] +pub struct ValidationError { + /// All failures found in a single validation pass. + pub failures: Vec, +} + +impl fmt::Display for ValidationError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + write!(f, "config validation failed ({} issue(s)):", self.failures.len())?; + for failure in &self.failures { + write!(f, "\n - {failure}")?; + } + Ok(()) + } +} + +impl std::error::Error for ValidationError {} + +impl From for ConfigError { + fn from(err: ValidationError) -> Self { + ConfigError::Validation(err.to_string()) + } +} + +/// Trait for types that can validate themselves after configuration loading. +/// +/// Implementors collect field-level failures rather than returning on the +/// first error, so operators see the full problem set in one pass. +pub trait ConfigValidator { + fn validate(&self) -> Result<(), ValidationError>; +} + +/// Convenience blanket for any serializable config type that implements +/// [`ConfigValidator`]. Callers typically invoke this on the output of +/// [`ConfigLoader::build`](crate::loader::ConfigLoader::build). +/// +/// ```no_run +/// # use mytheclipse_config::{ConfigLoader, ConfigValidator, ConfigValidatorExt}; +/// # use serde::Deserialize; +/// # #[derive(Debug, Deserialize)] +/// # struct Cfg { port: u16 } +/// # impl ConfigValidator for Cfg { +/// # fn validate(&self) -> Result<(), mytheclipse_config::ValidationError> { Ok(()) } +/// # } +/// let cfg: Cfg = ConfigLoader::new().build().unwrap(); +/// cfg.validate_config().unwrap(); +/// ``` +pub trait ConfigValidatorExt: ConfigValidator { + /// Validates `self`, returning `Ok(())` on success. + fn validate_config(&self) -> Result<(), ConfigError> { + self.validate().map_err(ConfigError::from) + } +} + +impl ConfigValidatorExt for T {} + +/// Validates that a string is a well-formed URL (http/https). +pub fn validate_url(path: &str, value: &str) -> Option { + if value.is_empty() { + return Some(ValidationFailure { + path: path.to_string(), + message: "url must not be empty".into(), + }); + } + // Minimal heuristic: scheme + host. We avoid pulling in a full URL crate + // to keep the dependency surface small. + let scheme_len = if value.starts_with("http://") { 7 } else if value.starts_with("https://") { 8 } else { + return Some(ValidationFailure { + path: path.to_string(), + message: format!("url must start with http:// or https:// (got {value:?})"), + }); + }; + let host = &value[scheme_len..]; + if host.is_empty() { + return Some(ValidationFailure { + path: path.to_string(), + message: format!("url has no host portion (got {value:?})"), + }); + } + None +} + +/// Validates that a port number is in the valid range (1–65535). +pub fn validate_port(path: &str, port: u16) -> Option { + // u16 already ranges 0–65535; exclude 0 (reserved/unspecified). + if port == 0 { + Some(ValidationFailure { + path: path.to_string(), + message: "port must be > 0".into(), + }) + } else { + None + } +} + +/// Validates that a string is non-empty. +pub fn validate_non_empty(path: &str, value: &str) -> Option { + if value.trim().is_empty() { + Some(ValidationFailure { + path: path.to_string(), + message: "value must not be empty".into(), + }) + } else { + None + } +} + +/// Validates that a numeric value falls within `[lo, hi]`. +pub fn validate_range(path: &str, value: T, lo: T, hi: T) -> Option +where + T: PartialOrd + fmt::Display + Copy, +{ + if value < lo || value > hi { + Some(ValidationFailure { + path: path.to_string(), + message: format!("value {value} is out of range [{lo}, {hi}]"), + }) + } else { + None + } +} + +/// Collects all failures from an iterator of `Option`. +pub fn collect_failures(opts: impl IntoIterator>) -> Result<(), ValidationError> { + let failures: Vec<_> = opts.into_iter().flatten().collect(); + if failures.is_empty() { + Ok(()) + } else { + Err(ValidationError { failures }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn url_validator_pass_and_fail() { + assert!(validate_url("db.url", "https://example.com").is_none()); + assert!(validate_url("db.url", "").is_some()); + assert!(validate_url("db.url", "ftp://bad").is_some()); + assert!(validate_url("db.url", "https://").is_some()); + } + + #[test] + fn port_validator_rejects_zero() { + assert!(validate_port("port", 0).is_some()); + assert!(validate_port("port", 1).is_none()); + assert!(validate_port("port", 65535).is_none()); + } + + #[test] + fn range_validator_bounds() { + assert!(validate_range("x", 5, 1, 10).is_none()); + assert!(validate_range("x", 10, 1, 10).is_none()); + assert!(validate_range("x", 0, 1, 10).is_some()); + assert!(validate_range("x", 11, 1, 10).is_some()); + } + + #[test] + fn collect_failures_aggregates_all() { + let opts = [validate_non_empty("a", ""), validate_non_empty("b", "ok"), validate_url("c.d", "bad://x")]; + let err = collect_failures(opts).unwrap_err(); + assert_eq!(err.failures.len(), 2); + assert_eq!(err.failures[0].path, "a"); + assert_eq!(err.failures[1].path, "c.d"); + } + + #[test] + fn collect_failures_ok_when_all_pass() { + let opts = [validate_url("a", "https://ok.com"), validate_port("b", 8080)]; + assert!(collect_failures(opts).is_ok()); + } + + #[test] + fn blanket_ext_wrappers_validator() { + struct Cfg; + impl ConfigValidator for Cfg { + fn validate(&self) -> Result<(), ValidationError> { + Err(ValidationError { + failures: vec![ValidationFailure { path: "x".into(), message: "bad".into() }], + }) + } + } + let c = Cfg; + assert!(c.validate_config().is_err()); + } +} diff --git a/crates/mytheclipse/src/lib.rs b/crates/mytheclipse/src/lib.rs index 35a3642..f818a65 100644 --- a/crates/mytheclipse/src/lib.rs +++ b/crates/mytheclipse/src/lib.rs @@ -55,11 +55,15 @@ pub mod cron; pub mod health; #[cfg(feature = "lifecycle")] pub mod leader; +#[cfg(feature = "lifecycle")] +pub mod lifecycle; #[cfg(feature = "observability")] pub mod metrics; #[cfg(feature = "observability")] pub mod panic_tracker; +#[cfg(feature = "observability")] +pub mod metrics_bridge; #[cfg(feature = "resiliency")] pub mod service_builder; @@ -106,10 +110,14 @@ pub use service_builder::ServiceBuilder; #[cfg(feature = "lifecycle")] pub use dlock::{DistributedLock, LockError, LockGuard, InProcLock}; +#[cfg(feature = "lifecycle")] +pub use lifecycle::AsyncLifecycleManager; #[cfg(feature = "observability")] pub use metrics::{MetricsCollector, MetricsSnapshot}; #[cfg(feature = "observability")] +pub use metrics_bridge::{MetricsBridge, MetricsHealthCheck}; +#[cfg(feature = "observability")] pub use panic_tracker::{PanicGuard, PanicInfo, PanicTracker}; /// Bootstraps the global [`EngineContext`]. diff --git a/crates/mytheclipse/src/lifecycle.rs b/crates/mytheclipse/src/lifecycle.rs new file mode 100644 index 0000000..48183cf --- /dev/null +++ b/crates/mytheclipse/src/lifecycle.rs @@ -0,0 +1,166 @@ +//! Async lifecycle manager composing shutdown, health checks, and periodic tasks. +//! +//! [`AsyncLifecycleManager`] ties together [`ShutdownManager`], [`HealthRegistry`], +//! and an optional periodic health-check ticker into a single orchestrator so +//! applications don't need to wire three separate primitives together. + +use std::sync::Arc; +use std::time::Duration; + +use crate::health::{HealthCheck, HealthRegistry, HealthStatus}; +use crate::shutdown::ShutdownManager; + +/// Coordinates graceful shutdown, health-check registration, and an optional +/// periodic health poll loop. +/// +/// Typical usage: +/// ```ignore +/// # tokio::runtime::Runtime::new().unwrap().block_on(async { +/// # use mytheclipse::AsyncLifecycleManager; +/// let mgr = AsyncLifecycleManager::new(); +/// mgr.register_health_check("db", my_db_check()); +/// let handle = mgr.start_health_loop(std::time::Duration::from_secs(30)); +/// mgr.await_shutdown(std::time::Duration::from_secs(10)).await; +/// ``` +#[derive(Clone)] +pub struct AsyncLifecycleManager { + shutdown: ShutdownManager, + health: Arc, +} + +impl AsyncLifecycleManager { + pub fn new() -> Self { + Self { + shutdown: ShutdownManager::new(), + health: Arc::new(HealthRegistry::new()), + } + } + + /// Returns a clone of the underlying shutdown manager. + pub fn shutdown(&self) -> &ShutdownManager { + &self.shutdown + } + + /// Returns a clone of the underlying health registry. + pub fn health(&self) -> &HealthRegistry { + &self.health + } + + /// Registers a named health check. + pub async fn register_health_check(&self, name: impl Into, check: impl HealthCheck + 'static) { + self.health.register(name, check).await; + } + + /// Runs all registered health checks once and returns their statuses. + pub async fn check_health(&self) -> Vec<(String, HealthStatus)> { + self.health.check_all().await + } + + /// Returns a shutdown signal for long-running tasks to observe. + pub fn shutdown_signal(&self) -> crate::shutdown::ShutdownSignal { + self.shutdown.handle() + } + + /// Starts a background task that polls health checks at `interval` and + /// emits tracing events. Returns a [`tokio::task::JoinHandle`] that can + /// be aborted on shutdown. + pub fn start_health_loop(&self, interval: Duration) -> tokio::task::JoinHandle<()> { + let health = self.health.clone(); + let signal = self.shutdown_signal(); + tokio::spawn(async move { + let mut ticker = tokio::time::interval(interval); + let mut sig = signal; + loop { + // Stop when shutdown is requested. + if sig.is_shutdown() { + tracing::info_span!("mytheclipse_health_loop", ); + return; + } + tokio::select! { + _ = sig.wait() => { + return; + } + _ = ticker.tick() => { + let results = health.check_all().await; + for (name, status) in &results { + match status { + HealthStatus::Ok => tracing::debug!(name, "health check ok"), + HealthStatus::Degraded => tracing::warn!(name, "health check degraded"), + HealthStatus::Unhealthy => tracing::error!(name, "health check unhealthy"), + } + } + if results.iter().any(|(_, s)| matches!(s, HealthStatus::Unhealthy)) { + tracing::error!("unhealthy component detected; requesting shutdown"); + return; + } + } + } + } + }) + } + + /// Waits for shutdown (OS signal or explicit `request()`) then drains all + /// registered tasks with a `grace` timeout per task. + pub async fn await_shutdown(&self, grace: Duration) { + self.shutdown.wait_for_shutdown().await; + self.shutdown.drain(grace).await; + } + + /// Requests shutdown programmatically (safe to call multiple times). + pub fn request_shutdown(&self) { + self.shutdown.request(); + } +} + +impl Default for AsyncLifecycleManager { + fn default() -> Self { + Self::new() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + struct AlwaysOk; + impl HealthCheck for AlwaysOk { + fn name(&self) -> &str { "always-ok" } + fn check(&self) -> std::pin::Pin + Send + '_>> { + Box::pin(async { HealthStatus::Ok }) + } + } + + struct AlwaysBad; + impl HealthCheck for AlwaysBad { + fn name(&self) -> &str { "always-bad" } + fn check(&self) -> std::pin::Pin + Send + '_>> { + Box::pin(async { HealthStatus::Unhealthy }) + } + } + + #[tokio::test] + async fn new_manager_has_no_checks() { + let mgr = AsyncLifecycleManager::new(); + let results = mgr.check_health().await; + assert!(results.is_empty()); + } + + #[tokio::test] + async fn registers_and_checks_health() { + let mgr = AsyncLifecycleManager::new(); + mgr.register_health_check("ok", AlwaysOk).await; + let results = mgr.check_health().await; + assert_eq!(results.len(), 1); + assert_eq!(results[0].0, "ok"); + assert_eq!(results[0].1, HealthStatus::Ok); + } + + #[tokio::test] + async fn shutdown_signal_fires_on_request() { + let mgr = AsyncLifecycleManager::new(); + let mut sig = mgr.shutdown_signal(); + assert!(!sig.is_shutdown()); + mgr.request_shutdown(); + assert!(sig.is_shutdown()); + } +} diff --git a/crates/mytheclipse/src/metrics_bridge.rs b/crates/mytheclipse/src/metrics_bridge.rs new file mode 100644 index 0000000..f2088bf --- /dev/null +++ b/crates/mytheclipse/src/metrics_bridge.rs @@ -0,0 +1,156 @@ +//! Bridges the metrics collector to health checks and tracing events. +//! +//! [`MetricsBridge`] ties [`crate::metrics::MetricsCollector`] to +//! [`crate::health::HealthCheck`], so a metrics-based health probe can report +//! `Degraded` when error counters rise or throughput drops, and optionally emit +//! tracing events so counters/gauges are visible in structured logs. + +use std::time::Duration; + +use crate::health::{HealthCheck, HealthStatus}; +use crate::metrics::MetricsCollector; + +/// A health check backed by a [`MetricsCollector`]: unhealthy if any registered +/// "error" counter is non-zero, degraded if any gauge is below a configured +/// threshold. +pub struct MetricsHealthCheck { + collector: MetricsCollector, +} + +impl MetricsHealthCheck { + pub fn new(collector: MetricsCollector) -> Self { + Self { collector } + } + + /// Returns unhealthy if the named counter is non-zero. + pub fn error_counter_exists(&self, name: &str) -> bool { + self.collector.snapshot().counters.contains_key(name) + } + + fn has_errors(&self) -> bool { + self.collector + .snapshot() + .counters + .values() + .any(|&v| v > 0) + } +} + +impl HealthCheck for MetricsHealthCheck { + fn name(&self) -> &str { + "metrics" + } + + fn check(&self) -> std::pin::Pin + Send + '_>> { + let has_errors = self.has_errors(); + Box::pin(async move { + if has_errors { + HealthStatus::Unhealthy + } else { + HealthStatus::Ok + } + }) + } +} + +/// Bridges a [`MetricsCollector`] to tracing: periodically emits the current +/// snapshot as tracing events so metrics are visible in structured logs. +pub struct MetricsBridge { + collector: MetricsCollector, +} + +impl MetricsBridge { + pub fn new(collector: MetricsCollector) -> Self { + Self { collector } + } + + /// Sends a one-shot tracing event with the current snapshot. + pub fn emit_now(&self) { + let snap = self.collector.snapshot(); + let mut counters: Vec<_> = snap.counters.into_iter().collect(); + counters.sort_by(|a, b| a.0.cmp(&b.0)); + let mut gauges: Vec<_> = snap.gauges.into_iter().collect(); + gauges.sort_by(|a, b| a.0.cmp(&b.0)); + + tracing::debug!( + task_count = snap.task_count, + active_threads = snap.active_threads, + queue_capacity_total = snap.queue_capacity_total, + queue_capacity_remaining = snap.queue_capacity_remaining, + "metrics snapshot" + ); + for (name, value) in &counters { + tracing::info!(name, value, "metric counter"); + } + for (name, value) in &gauges { + tracing::info!(name, value, "metric gauge"); + } + } + + /// Spawns a background task that calls [`emit_now`](Self::emit_now) every + /// `interval`. Returns a handle that can be aborted. + pub fn emit_periodic(self, interval: Duration) -> tokio::task::JoinHandle<()> { + tokio::spawn(async move { + let mut ticker = tokio::time::interval(interval); + loop { + ticker.tick().await; + self.emit_now(); + } + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::metrics::MetricsCollector; + use crate::shutdown::ShutdownManager; + + #[test] + fn metrics_health_ok_when_no_counters() { + let collector = MetricsCollector::new(); + let check = MetricsHealthCheck::new(collector); + // No counters set → no errors → Ok. + let fut = check.check(); + // Can't await in #[test]; use tokio test below instead. + drop(fut); + } + + #[tokio::test] + async fn metrics_health_unhealthy_when_errors_exist() { + let collector = MetricsCollector::new(); + collector.inc_counter("errors", 1); + let check = MetricsHealthCheck::new(collector); + let status = check.check().await; + assert_eq!(status, HealthStatus::Unhealthy); + } + + #[tokio::test] + async fn metrics_health_ok_when_no_errors() { + let collector = MetricsCollector::new(); + collector.set_gauge("load", 0.5); + let check = MetricsHealthCheck::new(collector); + let status = check.check().await; + assert_eq!(status, HealthStatus::Ok); + } + + #[tokio::test] + async fn bridge_emit_now_runs() { + let collector = MetricsCollector::new(); + collector.set_gauge("temp", 42.0); + let bridge = MetricsBridge::new(collector); + bridge.emit_now(); + } + + #[tokio::test] + async fn lifecycle_manager_with_metrics_bridge() { + let collector = MetricsCollector::new(); + collector.set_gauge("load", 0.1); + let mgr = crate::lifecycle::AsyncLifecycleManager::new(); + let bridge = MetricsBridge::new(collector); + let _handle = bridge.emit_periodic(Duration::from_millis(50)); + mgr.request_shutdown(); + // Should not hang — shutdown is immediate. + mgr.await_shutdown(Duration::from_secs(1)).await; + } +}