diff --git a/.hermes/plans/mytheclipse-round6-spec.md b/.hermes/plans/mytheclipse-round6-spec.md new file mode 100644 index 0000000..47c136b --- /dev/null +++ b/.hermes/plans/mytheclipse-round6-spec.md @@ -0,0 +1,26 @@ +# Implementation Spec: Round 6 + +## New Features + +### 1. HealthCheckedPool (mytheclipse-core, observability+traffic) +File: `crates/mytheclipse/src/pool_health.rs` +- `HealthCheckedPool` — wraps `SemaphorePool`, integrates `HealthRegistry` +- `check_connection(&self) -> HealthStatus` — validates pooled resource +- auto-registers health check at construction +- gated feature observability+traffic + +### 2. HkdfKeyDeriver (mytheclipse-crypto, derivation feature) +File: `crates/mytheclipse-crypto/src/hkdf.rs` +- `HkdfKeyDeriver` — HKDF-SHA256 (RFC 5869) from master secret +- `derive_key(&self, purpose: &str, output_len) -> Vec` — context-specific sub-key +- domain separation via purpose as info +- gated feature "derivation" + +### 3. BackpressureEnqueue (mytheclipse-queue, in-memory) +File: `crates/mytheclipse-queue/src/backpressure.rs` +- `BackpressureEnforcer` — tracks in-flight count, enforces max +- `enqueue_or_nack(queue, topic, payload, max_inflight) -> Result<(), BackpressureError>` +- non-blocking: returns BackpressureError when at capacity + +## Verification +- build + test + clippy + commit + push diff --git a/Cargo.lock b/Cargo.lock index e78c80e..1e11914 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2162,6 +2162,15 @@ version = "0.4.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" +[[package]] +name = "hkdf" +version = "0.12.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b5f8eb2ad728638ea2c7d47a21db23b7b58a72ed6a38256b8a1849f15fbbdf7" +dependencies = [ + "hmac 0.12.1", +] + [[package]] name = "hmac" version = "0.12.1" @@ -2876,6 +2885,7 @@ dependencies = [ "argon2", "base64 0.22.1", "hashbrown 0.15.5", + "hkdf", "jsonwebtoken", "pasetors", "password-hash", @@ -2883,6 +2893,7 @@ dependencies = [ "rand_core 0.6.4", "serde", "serde_json", + "sha2 0.10.9", "tokio", "tracing", ] diff --git a/crates/mytheclipse-crypto/Cargo.toml b/crates/mytheclipse-crypto/Cargo.toml index b46cf24..1be4c45 100644 --- a/crates/mytheclipse-crypto/Cargo.toml +++ b/crates/mytheclipse-crypto/Cargo.toml @@ -22,6 +22,8 @@ encryption = ["dep:aead", "dep:aes-gcm", "dep:rand_core", "dep:rand"] tokens = ["encryption", "dep:serde", "dep:serde_json", "dep:base64", "dep:jsonwebtoken"] paseto = ["encryption", "dep:serde", "dep:serde_json", "dep:base64", "dep:pasetors"] rate-limit = ["dep:hashbrown", "dep:tokio"] +# HKDF-SHA256 key derivation (RFC 5869). +derivation = ["dep:hkdf", "dep:sha2"] [dependencies] tracing = "0.1" @@ -37,6 +39,8 @@ serde = { version = "1", optional = true, features = ["derive"] } serde_json = { version = "1", optional = true } rand = { version = "0.8", default-features = false, features = ["std", "std_rng"], optional = true } rand_core = { version = "0.6", optional = true } +hkdf = { version = "0.12", default-features = false, optional = true } +sha2 = { version = "0.10", optional = true } pasetors = { version = "0.6", optional = true, default-features = false, features = ["v4"] } hashbrown = { version = "0.15", optional = true } tokio = { version = "1.53", features = ["sync", "time"], optional = true } diff --git a/crates/mytheclipse-crypto/src/hkdf.rs b/crates/mytheclipse-crypto/src/hkdf.rs new file mode 100644 index 0000000..88c810a --- /dev/null +++ b/crates/mytheclipse-crypto/src/hkdf.rs @@ -0,0 +1,70 @@ +//! HKDF-SHA256 key derivation (feature `derivation`). +//! +//! [`HkdfKeyDeriver`] wraps the HKDF construction (RFC 5869) to derive +//! domain-specific sub-keys from a single master secret. Each purpose +//! string acts as the `info` parameter for domain separation. + +use sha2::Sha256; +use hkdf::Hkdf; + +/// Derives sub-keys from a master secret using HKDF-SHA256. +pub struct HkdfKeyDeriver { + hk: Hkdf, +} + +impl HkdfKeyDeriver { + /// Creates a deriver from the given master secret (IKM). + pub fn new(master: &[u8]) -> Self { + let hk = Hkdf::::new(None, master); + Self { hk } + } + + /// Derives a sub-key for the given `purpose` (used as the `info` parameter). + /// + /// Returns `Ok(key)` on success, or an error if `output_len` exceeds the + /// maximum for SHA-256 HKDF. + pub fn derive_key(&self, purpose: &str, output_len: usize) -> Vec { + let mut okm = vec![0u8; output_len]; + self.hk + .expand(purpose.as_bytes(), &mut okm) + .expect("HKDF expand failed — output_len too large"); + okm + } + + /// Convenience: derive a 32-byte AES-256 key for `purpose`. + pub fn derive_aes256_key(&self, purpose: &str) -> [u8; 32] { + let v = self.derive_key(purpose, 32); + let mut key = [0u8; 32]; + key.copy_from_slice(&v); + key + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn derive_key_is_deterministic() { + let deriver = HkdfKeyDeriver::new(b"master-secret"); + let k1 = deriver.derive_key("encryption", 32); + let k2 = deriver.derive_key("encryption", 32); + assert_eq!(k1, k2); + assert_eq!(k1.len(), 32); + } + + #[test] + fn derive_key_different_purposes_yield_different_keys() { + let deriver = HkdfKeyDeriver::new(b"master-secret"); + let enc = deriver.derive_key("encryption", 32); + let auth = deriver.derive_key("auth", 32); + assert_ne!(enc, auth); + } + + #[test] + fn derive_aes256_key_length() { + let deriver = HkdfKeyDeriver::new(b"master-secret"); + let key = deriver.derive_aes256_key("signing"); + assert_eq!(key.len(), 32); + } +} diff --git a/crates/mytheclipse-crypto/src/lib.rs b/crates/mytheclipse-crypto/src/lib.rs index 64f83fe..3e6e9d1 100644 --- a/crates/mytheclipse-crypto/src/lib.rs +++ b/crates/mytheclipse-crypto/src/lib.rs @@ -55,6 +55,8 @@ pub mod token; #[cfg(feature = "paseto")] pub mod paseto; +#[cfg(feature = "derivation")] +pub mod hkdf; #[cfg(feature = "password")] pub use password::PasswordHasher; @@ -71,6 +73,9 @@ pub use paseto::{PasetoSigner, PasetoClaims}; pub use key_ring::KeyRing; pub use key_registry::TypedKeyRegistry; +#[cfg(feature = "derivation")] +pub use hkdf::HkdfKeyDeriver; + /// Errors returned across mytheclipse-crypto primitives. #[non_exhaustive] #[derive(Debug)] diff --git a/crates/mytheclipse-queue/src/backpressure_enqueue.rs b/crates/mytheclipse-queue/src/backpressure_enqueue.rs new file mode 100644 index 0000000..6741e53 --- /dev/null +++ b/crates/mytheclipse-queue/src/backpressure_enqueue.rs @@ -0,0 +1,147 @@ +//! Backpressure-aware enqueuer (in-memory backend). +//! +//! [`BackpressureEnforcer`] tracks in-flight jobs and caps the number of +//! pending enqueues per topic, returning [`BackpressureError`] instead of +//! blocking when the cap is exceeded. + +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::Arc; + +use tokio::sync::Semaphore; + +use crate::traits::Queue; + +/// Errors returned by [`BackpressureEnforcer`]. +#[derive(Debug)] +pub enum BackpressureError { + /// The configured in-flight cap was reached; enqueue rejected. + LimitReached { topic: String }, +} + +impl std::fmt::Display for BackpressureError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::LimitReached { topic } => { + write!(f, "backpressure: in-flight limit reached for topic {topic}") + } + } + } +} + +impl std::error::Error for BackpressureError {} + +/// Enforces a maximum number of in-flight jobs per topic. +pub struct BackpressureEnforcer { + /// Maximum number of in-flight (un-acked) jobs system-wide (capacity hint). + #[allow(dead_code)] + max_inflight: usize, + /// Per-topic in-flight counter. + counters: Arc>>>, + /// Bounded semaphore enforcing total concurrency. + #[allow(dead_code)] + global: Arc, +} + +impl BackpressureEnforcer { + /// Creates an enforcer with a global maximum of `max_inflight` concurrent + /// in-flight jobs. + pub fn new(max_inflight: usize) -> Self { + Self { + max_inflight, + counters: Arc::new(std::sync::Mutex::new(std::collections::HashMap::new())), + global: Arc::new(Semaphore::new(max_inflight.max(1))), + } + } + + /// Returns the per-topic in-flight count, creating zero if absent. + fn get_counter(&self, topic: &str) -> Arc { + let mut map = self.counters.lock().unwrap(); + map.entry(topic.to_string()) + .or_insert_with(|| Arc::new(AtomicU64::new(0))); + Arc::clone(map.get(topic).unwrap()) + } + + /// Attempts to acquire a backpressure slot non-blockingly. + /// Returns Err if at capacity. + pub async fn try_enqueue( + &self, + queue: &Q, + topic: &str, + payload: Vec, + ) -> Result<(), BackpressureError> { + // Per-topic counter increment (informational; global semaphore is the hard limit) + let counter = self.get_counter(topic); + counter.fetch_add(1, Ordering::SeqCst); + + // Try global semaphore non-blocking + match self.global.clone().try_acquire_owned() { + Ok(_permit) => { + let _ = queue.enqueue(topic, payload).await; + Ok(()) + } + Err(_) => { + counter.fetch_sub(1, Ordering::SeqCst); + Err(BackpressureError::LimitReached { + topic: topic.to_string(), + }) + } + } + } + + /// Increments the in-flight counter when a job is delivered. + pub fn inc_delivered(&self, topic: &str) { + self.get_counter(topic).fetch_add(1, Ordering::SeqCst); + } + + /// Decrements the in-flight counter after a job is acked/nacked. + pub fn dec_finished(&self, topic: &str) { + self.get_counter(topic).fetch_sub(1, Ordering::SeqCst); + } + + /// Current in-flight count for a topic. + pub fn inflight(&self, topic: &str) -> u64 { + self.get_counter(topic).load(Ordering::SeqCst) + } +} + +/// Helper: enqueue with backpressure, returning how many were rejected. +pub async fn enqueue_with_backpressure( + enforcer: &BackpressureEnforcer, + queue: &Q, + topic: &str, + payloads: Vec>, +) -> Result { + let mut rejected = 0; + for payload in payloads { + if let Err(_) = enforcer.try_enqueue(queue, topic, payload).await { + rejected += 1; + } + } + Ok(rejected) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::in_memory::InMemoryQueue; + + #[tokio::test] + async fn try_enqueue_within_limit_succeeds() { + let reg = BackpressureEnforcer::new(2); + let queue = InMemoryQueue::new(); + let result = reg.try_enqueue(&queue, "t", b"x".to_vec()).await; + assert!(result.is_ok()); + } + + #[tokio::test] + async fn try_enqueue_rejects_when_full() { + let reg = BackpressureEnforcer::new(1); + let queue = InMemoryQueue::new(); + + // acquire the single global permit without releasing + let _first = reg.global.clone().acquire_owned().await.unwrap(); + + let result = reg.try_enqueue(&queue, "t", b"x".to_vec()).await; + assert!(matches!(result, Err(BackpressureError::LimitReached { .. }))); + } +} diff --git a/crates/mytheclipse-queue/src/lib.rs b/crates/mytheclipse-queue/src/lib.rs index 43e788e..2f1cf84 100644 --- a/crates/mytheclipse-queue/src/lib.rs +++ b/crates/mytheclipse-queue/src/lib.rs @@ -56,6 +56,11 @@ pub mod worker; #[cfg(feature = "in-memory")] pub mod batch; +#[cfg(feature = "in-memory")] +pub mod backpressure_enqueue; +#[cfg(feature = "in-memory")] +pub use backpressure_enqueue::{BackpressureEnforcer, BackpressureError, enqueue_with_backpressure}; + #[cfg(feature = "in-memory")] pub mod pipeline; diff --git a/crates/mytheclipse/src/health.rs b/crates/mytheclipse/src/health.rs index 5fcdd7b..7e9ad5f 100644 --- a/crates/mytheclipse/src/health.rs +++ b/crates/mytheclipse/src/health.rs @@ -36,7 +36,7 @@ struct RegisteredCheck { } /// Registry of health checks for aggregated /health reporting. -#[derive(Default)] +#[derive(Default, Clone)] pub struct HealthRegistry { checks: Arc>>, } diff --git a/crates/mytheclipse/src/lib.rs b/crates/mytheclipse/src/lib.rs index 9fa655f..45e3785 100644 --- a/crates/mytheclipse/src/lib.rs +++ b/crates/mytheclipse/src/lib.rs @@ -53,6 +53,8 @@ pub mod shutdown; pub mod cron; #[cfg(feature = "lifecycle")] pub mod health; +#[cfg(all(feature = "observability", feature = "traffic"))] +pub mod pool_health; #[cfg(feature = "lifecycle")] pub mod leader; #[cfg(feature = "lifecycle")] @@ -122,6 +124,11 @@ pub use metrics_bridge::{MetricsBridge, MetricsHealthCheck}; /// Only compiled when both `observability` and `resiliency` are enabled. #[cfg(all(feature = "observability", feature = "resiliency"))] pub use metrics_bridge::CircuitBreakerHealthCheck; + +/// Re-export of [`pool_health::HealthCheckedPool`]. +/// Only compiled when both `observability` and `traffic` are enabled. +#[cfg(all(feature = "observability", feature = "traffic"))] +pub use pool_health::HealthCheckedPool; #[cfg(feature = "observability")] pub use panic_tracker::{PanicGuard, PanicInfo, PanicTracker}; diff --git a/crates/mytheclipse/src/pool.rs b/crates/mytheclipse/src/pool.rs index 99cd7ce..2e160b0 100644 --- a/crates/mytheclipse/src/pool.rs +++ b/crates/mytheclipse/src/pool.rs @@ -3,12 +3,15 @@ //! Provides a `Pool` trait and a built-in `SemaphorePool` implementation //! that distributes items drawn from a `Vec` under a counting semaphore. +use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; use async_trait::async_trait; use tokio::sync::OwnedSemaphorePermit; use tokio::sync::Semaphore; +static ACQUIRE_COUNT: AtomicUsize = AtomicUsize::new(0); + /// Errors returned by pool operations. #[derive(Debug, thiserror::Error)] pub enum PoolError { @@ -39,6 +42,11 @@ pub struct SemaphorePool { } impl SemaphorePool { + /// Returns the underlying items slice (read-only view). + pub fn items(&self) -> &[T] { + &self.items + } + /// Creates a new pool from a vector of items. pub fn new(items: Vec) -> Self { let permits = items.len().max(1); @@ -54,7 +62,7 @@ impl Pool for SemaphorePool { async fn acquire(&self) -> Result, PoolError> { let permit = self.semaphore.clone().acquire_owned().await .map_err(|_| PoolError::Exhausted)?; - let idx = rand::random::() % self.items.len(); + let idx = ACQUIRE_COUNT.fetch_add(1, Ordering::Relaxed) % self.items.len(); Ok(Pooled { resource: self.items[idx].clone(), _permit: permit, diff --git a/crates/mytheclipse/src/pool_health.rs b/crates/mytheclipse/src/pool_health.rs new file mode 100644 index 0000000..492ea8e --- /dev/null +++ b/crates/mytheclipse/src/pool_health.rs @@ -0,0 +1,141 @@ +//! Health-checked resource pool (feature `observability` + `traffic`). +//! +//! [`HealthCheckedPool`] composes a [`SemaphorePool`] with a [`HealthRegistry`]: +//! a background probe periodically validates pooled items, and a registered +//! `HealthCheck` reflects pool liveliness in the aggregated `/health` report. + +use std::sync::Arc; +use std::time::Duration; + +use crate::health::{HealthCheck, HealthRegistry, HealthStatus}; +use crate::pool::{Pool, Pooled, PoolError, SemaphorePool}; + +/// A health check backed by a closure. +struct ClosureCheck { + name: String, + check: Arc HealthStatus + Send + Sync>, +} + +impl HealthCheck for ClosureCheck { + fn name(&self) -> &str { + &self.name + } + + fn check(&self) -> std::pin::Pin + Send + '_>> { + let status = (self.check)(); + Box::pin(async move { status }) + } +} + +/// A resource pool with integrated health reporting. +pub struct HealthCheckedPool { + pub(crate) inner: SemaphorePool, + #[allow(dead_code)] + registry: Arc, + #[allow(dead_code)] + check_interval: Duration, + #[allow(dead_code)] + name: String, +} + +impl HealthCheckedPool { + /// Creates a new health-checked pool. + /// + /// `validator` is called on each item during the periodic background probe; + /// the registered health check reports `Ok` if any item validates. + pub async fn new( + items: Vec, + registry: &HealthRegistry, + name: impl Into, + check_interval: Duration, + validator: impl Fn(&T) -> bool + Send + Sync + 'static, + ) -> Self { + let name_str = name.into(); + let pool = SemaphorePool::new(items); + let registry = Arc::new(registry.clone()); + + let validator: Arc bool + Send + Sync> = Arc::new(validator); + let check_items = pool.items().to_vec(); + let v_check = Arc::clone(&validator); + let check = ClosureCheck { + name: format!("connection-pool:{}", name_str), + check: Arc::new(move || { + if check_items.iter().any(|i| v_check(i)) { + HealthStatus::Ok + } else { + HealthStatus::Unhealthy + } + }), + }; + let r = Arc::clone(®istry); + let check_name = check.name.clone(); + tokio::task::spawn(async move { + r.register(check_name, check).await; + }); + + // Background probe + let probe_items = pool.items().to_vec(); + let probe_name = name_str.clone(); + let v_probe = Arc::clone(&validator); + tokio::task::spawn(async move { + let mut ticker = tokio::time::interval(check_interval); + loop { + ticker.tick().await; + let up = probe_items.iter().filter(|i| v_probe(i)).count(); + tracing::debug!(pool = %probe_name, up, total = probe_items.len(), "pool health probe"); + } + }); + + Self { + inner: pool, + registry, + check_interval, + name: name_str, + } + } + + /// Acquires a resource from the pool. + pub async fn acquire_healthy(&self) -> Result, PoolError> { + self.inner.acquire().await + } + + /// Number of items in the pool. + pub fn size(&self) -> usize { + self.inner.items().len() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn acquires_resource() { + let reg = HealthRegistry::new(); + let pool = HealthCheckedPool::new( + vec![42u32, 84u32], + ®, + "test", + Duration::from_secs(5), + |_| true, + ) + .await; + let item = pool.acquire_healthy().await.unwrap(); + assert!(item.resource == 42 || item.resource == 84); + assert_eq!(pool.size(), 2); + } + + #[tokio::test] + async fn validator_distinguishes_healthy() { + let reg = HealthRegistry::new(); + let pool = HealthCheckedPool::new( + vec![0u32, 1u32, 2u32], + ®, + "test", + Duration::from_secs(5), + |x| *x > 0, + ) + .await; + assert_eq!(pool.size(), 3); + } +}