From 1dfc6d68657e5b002896f400d744b052800c9285 Mon Sep 17 00:00:00 2001 From: asepharyana Date: Sun, 30 Aug 2026 00:33:08 +0700 Subject: [PATCH] =?UTF-8?q?fix(ci):=20restore=20full=20CI=20green=20?= =?UTF-8?q?=E2=80=94=20test-matrix,=20clippy,=20rustfmt,=20and=20rustdoc?= =?UTF-8?q?=20gates?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Root cause of the failing CI run was that examples/tests referencing feature-gated items were auto-detected (no required-features), so `--all-targets` compiled them under feature combinations where those modules didn't exist. Fixes: - mytheclipse Cargo.toml: declare the `high_level` example and `race_stress` integration test with required-features = ["full"]; `cargo build --all-targets` now skips them when full is off. This clears the whole test-matrix (workspace default, all-features, and every single-feature config) which all failed on E0432/E0433. - lib.rs: auto_metrics_service depends on service_builder, so regate it behind all(observability, resiliency) instead of observability alone (observability-only build compiled the module without resiliency). - mytheclipse-tracing: gate `pub mod fmt` behind any tracing-subscriber-providing feature so --no-default-features compiles. - mytheclipse-queue: gate `pub mod worker` behind in-memory (worker.rs requires tokio, only provided by in-memory). - clippy -D warnings fixes: deprecated base64 0.22 free fns -> Engine (paseto), unused key field, needless mut (service_builder), unused import/dead var/missing is_empty (bg_join), dead is_expired (dlock), while-let-iterator->for (parallel_map), type_complexity (shutdown_guard), MutexGuard held across await (middleware, now clones Arc'd layers), if-let-Err->is_err (queue), unused CliBuilder fields now wired into clap. - rustdoc -D warnings: resolve retry/MetricsBridge/CircuitBreaker/KeyRing intra-doc links and fix the unparseable lifecycle.rs code fence. - cargo fmt --all to satisfy the Rustfmt gate. --- crates/mytheclipse-cache/src/auto_refresh.rs | 7 +- crates/mytheclipse-cache/src/metrics.rs | 6 +- crates/mytheclipse-cli/src/builder.rs | 15 ++- crates/mytheclipse-config/src/lib.rs | 5 +- crates/mytheclipse-config/src/validate.rs | 32 +++++- crates/mytheclipse-crypto/src/hkdf.rs | 2 +- crates/mytheclipse-crypto/src/key_registry.rs | 11 +- crates/mytheclipse-crypto/src/lib.rs | 10 +- crates/mytheclipse-crypto/src/paseto.rs | 27 +++-- crates/mytheclipse-http/src/lib.rs | 2 +- crates/mytheclipse-http/src/metrics_http.rs | 6 +- .../mytheclipse-http/src/resilient_client.rs | 15 +-- .../src/server/axum_server.rs | 10 +- .../src/backpressure_enqueue.rs | 7 +- crates/mytheclipse-queue/src/batch.rs | 22 +++- crates/mytheclipse-queue/src/in_memory.rs | 16 ++- crates/mytheclipse-queue/src/lib.rs | 15 ++- crates/mytheclipse-queue/src/pipeline.rs | 14 +-- crates/mytheclipse-queue/src/rate_limited.rs | 14 ++- crates/mytheclipse-queue/src/traits.rs | 2 +- crates/mytheclipse-queue/src/worker.rs | 15 ++- .../src/worker_rate_limited.rs | 10 +- .../tests/rate_limited_stress.rs | 8 +- crates/mytheclipse-storage/src/local.rs | 10 +- crates/mytheclipse-tracing/src/fmt.rs | 4 +- crates/mytheclipse-tracing/src/lib.rs | 12 ++ crates/mytheclipse/Cargo.toml | 10 ++ crates/mytheclipse/benches/primitives.rs | 14 ++- crates/mytheclipse/examples/high_level.rs | 31 +++--- crates/mytheclipse/examples/scaling_demo.rs | 17 ++- crates/mytheclipse/src/aggregate_error.rs | 13 ++- .../mytheclipse/src/auto_metrics_service.rs | 33 ++++-- crates/mytheclipse/src/bg_join.rs | 15 ++- crates/mytheclipse/src/compute.rs | 76 ++++++------- crates/mytheclipse/src/dlock.rs | 54 ++++++--- crates/mytheclipse/src/health.rs | 4 +- crates/mytheclipse/src/lib.rs | 82 +++++++------- crates/mytheclipse/src/lifecycle.rs | 30 +++-- crates/mytheclipse/src/metrics_bridge.rs | 20 ++-- crates/mytheclipse/src/middleware.rs | 20 ++-- crates/mytheclipse/src/parallel_map.rs | 58 ++++------ crates/mytheclipse/src/pool.rs | 12 +- crates/mytheclipse/src/pool_health.rs | 6 +- crates/mytheclipse/src/retry.rs | 17 ++- crates/mytheclipse/src/retry_ext.rs | 16 ++- crates/mytheclipse/src/service_builder.rs | 105 ++++++++++++------ crates/mytheclipse/src/shutdown_guard.rs | 6 +- crates/mytheclipse/tests/race_stress.rs | 6 +- 48 files changed, 574 insertions(+), 368 deletions(-) diff --git a/crates/mytheclipse-cache/src/auto_refresh.rs b/crates/mytheclipse-cache/src/auto_refresh.rs index fb887a6..290aef0 100644 --- a/crates/mytheclipse-cache/src/auto_refresh.rs +++ b/crates/mytheclipse-cache/src/auto_refresh.rs @@ -71,7 +71,12 @@ where } /// Sets a value in the underlying cache. - pub async fn set(&self, key: &str, value: Vec, ttl: Option) -> Result<(), CacheError> { + pub async fn set( + &self, + key: &str, + value: Vec, + ttl: Option, + ) -> Result<(), CacheError> { self.inner.set(key, value, ttl).await } diff --git a/crates/mytheclipse-cache/src/metrics.rs b/crates/mytheclipse-cache/src/metrics.rs index 5fb878c..c763068 100644 --- a/crates/mytheclipse-cache/src/metrics.rs +++ b/crates/mytheclipse-cache/src/metrics.rs @@ -54,6 +54,10 @@ pub struct CacheSnapshot { impl CacheSnapshot { pub fn hit_rate(&self) -> f64 { let total = self.hits + self.misses; - if total == 0 { 0.0 } else { self.hits as f64 / total as f64 } + if total == 0 { + 0.0 + } else { + self.hits as f64 / total as f64 + } } } diff --git a/crates/mytheclipse-cli/src/builder.rs b/crates/mytheclipse-cli/src/builder.rs index 9b7d25b..9ae3d96 100644 --- a/crates/mytheclipse-cli/src/builder.rs +++ b/crates/mytheclipse-cli/src/builder.rs @@ -1,6 +1,6 @@ //! Clap-based CLI builder implementation. -use clap::{Parser, Subcommand as ClapSubcommand}; +use clap::{CommandFactory, FromArgMatches, Parser, Subcommand as ClapSubcommand}; /// A mytheclipse CLI application. #[derive(Parser, Debug)] @@ -52,6 +52,17 @@ impl CliBuilder { } pub fn build(self) -> CliApp { - CliApp::parse() + // Apply the configured name/about to the derived clap Command so the + // builder's fields are honored in the rendered help/usage. + let Self { name, about } = self; + // clap's `Str`/`StyledStr` only accept 'static references, so leak + // the owned strings (build(self) consumes self once, so a single, + // process-lifetime leak is acceptable). + let name: &'static str = String::leak(name); + let about: &'static str = String::leak(about); + let cmd = ::command() + .name(name) + .about(about); + CliApp::from_arg_matches(&cmd.get_matches()).unwrap_or_else(|e| e.exit()) } } diff --git a/crates/mytheclipse-config/src/lib.rs b/crates/mytheclipse-config/src/lib.rs index a526f66..1f9d8d2 100644 --- a/crates/mytheclipse-config/src/lib.rs +++ b/crates/mytheclipse-config/src/lib.rs @@ -51,9 +51,8 @@ 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, + collect_failures, validate_non_empty, validate_port, validate_range, validate_url, + ConfigValidator, ConfigValidatorExt, ValidationError, ValidationFailure, }; #[cfg(feature = "hot-reload")] diff --git a/crates/mytheclipse-config/src/validate.rs b/crates/mytheclipse-config/src/validate.rs index 09c01aa..3773d78 100644 --- a/crates/mytheclipse-config/src/validate.rs +++ b/crates/mytheclipse-config/src/validate.rs @@ -33,7 +33,11 @@ pub struct ValidationError { impl fmt::Display for ValidationError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - write!(f, "config validation failed ({} issue(s)):", self.failures.len())?; + write!( + f, + "config validation failed ({} issue(s)):", + self.failures.len() + )?; for failure in &self.failures { write!(f, "\n - {failure}")?; } @@ -91,7 +95,11 @@ pub fn validate_url(path: &str, value: &str) -> Option { } // 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 { + 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:?})"), @@ -148,7 +156,9 @@ where } /// Collects all failures from an iterator of `Option`. -pub fn collect_failures(opts: impl IntoIterator>) -> Result<(), ValidationError> { +pub fn collect_failures( + opts: impl IntoIterator>, +) -> Result<(), ValidationError> { let failures: Vec<_> = opts.into_iter().flatten().collect(); if failures.is_empty() { Ok(()) @@ -186,7 +196,11 @@ mod tests { #[test] fn collect_failures_aggregates_all() { - let opts = [validate_non_empty("a", ""), validate_non_empty("b", "ok"), validate_url("c.d", "bad://x")]; + 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"); @@ -195,7 +209,10 @@ mod tests { #[test] fn collect_failures_ok_when_all_pass() { - let opts = [validate_url("a", "https://ok.com"), validate_port("b", 8080)]; + let opts = [ + validate_url("a", "https://ok.com"), + validate_port("b", 8080), + ]; assert!(collect_failures(opts).is_ok()); } @@ -205,7 +222,10 @@ mod tests { impl ConfigValidator for Cfg { fn validate(&self) -> Result<(), ValidationError> { Err(ValidationError { - failures: vec![ValidationFailure { path: "x".into(), message: "bad".into() }], + failures: vec![ValidationFailure { + path: "x".into(), + message: "bad".into(), + }], }) } } diff --git a/crates/mytheclipse-crypto/src/hkdf.rs b/crates/mytheclipse-crypto/src/hkdf.rs index 88c810a..a5c1cf7 100644 --- a/crates/mytheclipse-crypto/src/hkdf.rs +++ b/crates/mytheclipse-crypto/src/hkdf.rs @@ -4,8 +4,8 @@ //! 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; +use sha2::Sha256; /// Derives sub-keys from a master secret using HKDF-SHA256. pub struct HkdfKeyDeriver { diff --git a/crates/mytheclipse-crypto/src/key_registry.rs b/crates/mytheclipse-crypto/src/key_registry.rs index 2c3af8b..1976845 100644 --- a/crates/mytheclipse-crypto/src/key_registry.rs +++ b/crates/mytheclipse-crypto/src/key_registry.rs @@ -1,6 +1,6 @@ //! Typed key registry with ID-based lookup (feature `password`). //! -//! [`TypedKeyRegistry`] extends [`KeyRing`] semantics: instead of a single +//! [`TypedKeyRegistry`] extends `KeyRing` semantics: instead of a single //! current+previous sequence, it maintains a map of named keys keyed by an ID, //! with one designated "current" ID. This is useful when keys are rotated by ID //! (e.g. JWT `kid` header) and you need to look up a verification key by ID @@ -20,7 +20,10 @@ pub struct TypedKeyRegistry { impl TypedKeyRegistry { /// Creates an empty registry (no current key). pub fn new() -> Self { - Self { keys: HashMap::new(), current_id: None } + Self { + keys: HashMap::new(), + current_id: None, + } } /// Registers a key under `id`, making it the current key. @@ -37,9 +40,7 @@ impl TypedKeyRegistry { /// Returns the current key, if any. pub fn current(&self) -> Option<&T> { - self.current_id - .as_ref() - .and_then(|id| self.keys.get(id)) + self.current_id.as_ref().and_then(|id| self.keys.get(id)) } /// Returns the ID of the current key. diff --git a/crates/mytheclipse-crypto/src/lib.rs b/crates/mytheclipse-crypto/src/lib.rs index 3e6e9d1..2525440 100644 --- a/crates/mytheclipse-crypto/src/lib.rs +++ b/crates/mytheclipse-crypto/src/lib.rs @@ -41,8 +41,8 @@ //! assert_eq!(claims["sub"], "u1"); //! ``` -pub mod key_ring; pub mod key_registry; +pub mod key_ring; #[cfg(feature = "password")] pub mod password; @@ -53,10 +53,10 @@ pub mod encryption; #[cfg(feature = "tokens")] pub mod token; -#[cfg(feature = "paseto")] -pub mod paseto; #[cfg(feature = "derivation")] pub mod hkdf; +#[cfg(feature = "paseto")] +pub mod paseto; #[cfg(feature = "password")] pub use password::PasswordHasher; @@ -68,10 +68,10 @@ pub use encryption::{AeadError, Encryptor}; pub use token::{Claims, TokenError, TokenSigner}; #[cfg(feature = "paseto")] -pub use paseto::{PasetoSigner, PasetoClaims}; +pub use paseto::{PasetoClaims, PasetoSigner}; -pub use key_ring::KeyRing; pub use key_registry::TypedKeyRegistry; +pub use key_ring::KeyRing; #[cfg(feature = "derivation")] pub use hkdf::HkdfKeyDeriver; diff --git a/crates/mytheclipse-crypto/src/paseto.rs b/crates/mytheclipse-crypto/src/paseto.rs index 71b63b3..b62b38b 100644 --- a/crates/mytheclipse-crypto/src/paseto.rs +++ b/crates/mytheclipse-crypto/src/paseto.rs @@ -5,6 +5,8 @@ use std::time::{Duration, SystemTime}; +use base64::engine::general_purpose::STANDARD; +use base64::Engine; use serde::{Deserialize, Serialize}; /// Errors returned by PASETO operations. @@ -58,26 +60,26 @@ impl PasetoClaims { /// /// This is a stub implementation. For production use with `pasetors` 0.6, /// the token format follows the PASETO v4.local specification. -pub struct PasetoSigner { - key: Vec, -} +pub struct PasetoSigner {} impl PasetoSigner { /// Creates a new signer with the given 32-byte key. pub fn new(key: &[u8]) -> Result { if key.len() != 32 { - return Err(PasetoError::Sign("key must be 32 bytes for v4-local".to_string())); + return Err(PasetoError::Sign( + "key must be 32 bytes for v4-local".to_string(), + )); } - Ok(Self { key: key.to_vec() }) + Ok(Self {}) } /// Signs claims into a PASETO v4.local token string. pub fn sign(&self, claims: &PasetoClaims) -> Result { - let payload = serde_json::to_string(claims) - .map_err(|e| PasetoError::Sign(e.to_string()))?; + let payload = + serde_json::to_string(claims).map_err(|e| PasetoError::Sign(e.to_string()))?; let nonce = rand::random::<[u8; 24]>(); - let nonce_b64 = base64::encode(&nonce); - let payload_b64 = base64::encode(payload); + let nonce_b64 = STANDARD.encode(nonce); + let payload_b64 = STANDARD.encode(payload.as_bytes()); Ok(format!("v4.local.{nonce_b64}.{payload_b64}")) } @@ -88,10 +90,11 @@ impl PasetoSigner { return Err(PasetoError::InvalidToken); } - let payload_bytes = base64::decode(parts[3]) - .map_err(|_| PasetoError::InvalidToken)?; - let claims: PasetoClaims = serde_json::from_slice(&payload_bytes) + let payload_bytes = STANDARD + .decode(parts[3]) .map_err(|_| PasetoError::InvalidToken)?; + let claims: PasetoClaims = + serde_json::from_slice(&payload_bytes).map_err(|_| PasetoError::InvalidToken)?; let now = SystemTime::now() .duration_since(SystemTime::UNIX_EPOCH) diff --git a/crates/mytheclipse-http/src/lib.rs b/crates/mytheclipse-http/src/lib.rs index 71997e4..1a7970a 100644 --- a/crates/mytheclipse-http/src/lib.rs +++ b/crates/mytheclipse-http/src/lib.rs @@ -12,7 +12,7 @@ #[cfg(feature = "resilience")] pub mod resilient_client; #[cfg(feature = "resilience")] -pub use resilient_client::{ResilientHttpClient, ResilientClientConfig}; +pub use resilient_client::{ResilientClientConfig, ResilientHttpClient}; #[cfg(feature = "client")] pub mod client; diff --git a/crates/mytheclipse-http/src/metrics_http.rs b/crates/mytheclipse-http/src/metrics_http.rs index 7409dcd..f8f6d76 100644 --- a/crates/mytheclipse-http/src/metrics_http.rs +++ b/crates/mytheclipse-http/src/metrics_http.rs @@ -28,11 +28,7 @@ async fn metrics_handler( .status(200) .header("content-type", "text/plain; version=0.0.4") .body(axum::body::Body::from(body)) - .unwrap_or_else(|_| { - axum::response::Response::new(axum::body::Body::from( - "internal error", - )) - }) + .unwrap_or_else(|_| axum::response::Response::new(axum::body::Body::from("internal error"))) } #[cfg(test)] diff --git a/crates/mytheclipse-http/src/resilient_client.rs b/crates/mytheclipse-http/src/resilient_client.rs index d8887a4..1becc6f 100644 --- a/crates/mytheclipse-http/src/resilient_client.rs +++ b/crates/mytheclipse-http/src/resilient_client.rs @@ -77,25 +77,22 @@ impl ResilientHttpClient { /// Sends a pre-built `RequestBuilder` through the resiliency pipeline. /// Returns the response bytes on success. - pub async fn send( - &self, - req: RequestBuilder, - ) -> Result, RunError> { + pub async fn send(&self, req: RequestBuilder) -> Result, RunError> { let span = tracing::info_span!("resilient_http_send"); let op = move || { let req = req.try_clone().unwrap(); let fut: Pin, HttpError>> + Send>> = Box::pin(async move { - let resp = req.send().instrument(tracing::trace_span!("http_send")).await?; + let resp = req + .send() + .instrument(tracing::trace_span!("http_send")) + .await?; let bytes = resp.bytes().await?; Ok::, HttpError>(bytes.to_vec()) }); fut }; - self.builder - .run(op) - .instrument(span) - .await + self.builder.run(op).instrument(span).await } /// Convenience: GET `url`, returning response bytes. diff --git a/crates/mytheclipse-http/src/server/axum_server.rs b/crates/mytheclipse-http/src/server/axum_server.rs index 369dfdf..e0694c6 100644 --- a/crates/mytheclipse-http/src/server/axum_server.rs +++ b/crates/mytheclipse-http/src/server/axum_server.rs @@ -1,9 +1,6 @@ //! Axum-based HTTP server with health endpoint. -use axum::{ - routing::get, - Router, -}; +use axum::{routing::get, Router}; use std::net::SocketAddr; /// A pre-configured HTTP server with health check and metrics endpoints. @@ -19,10 +16,7 @@ impl HttpServer { .route("/health", get(|| async { "OK" })) .route("/", get(|| async { "mytheclipse-http" })); - Self { - app: router, - addr, - } + Self { app: router, addr } } /// Adds a custom route with a GET handler. diff --git a/crates/mytheclipse-queue/src/backpressure_enqueue.rs b/crates/mytheclipse-queue/src/backpressure_enqueue.rs index 6741e53..a88556f 100644 --- a/crates/mytheclipse-queue/src/backpressure_enqueue.rs +++ b/crates/mytheclipse-queue/src/backpressure_enqueue.rs @@ -113,7 +113,7 @@ pub async fn enqueue_with_backpressure( ) -> Result { let mut rejected = 0; for payload in payloads { - if let Err(_) = enforcer.try_enqueue(queue, topic, payload).await { + if enforcer.try_enqueue(queue, topic, payload).await.is_err() { rejected += 1; } } @@ -142,6 +142,9 @@ mod tests { 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 { .. }))); + assert!(matches!( + result, + Err(BackpressureError::LimitReached { .. }) + )); } } diff --git a/crates/mytheclipse-queue/src/batch.rs b/crates/mytheclipse-queue/src/batch.rs index 93a247a..27f0b7b 100644 --- a/crates/mytheclipse-queue/src/batch.rs +++ b/crates/mytheclipse-queue/src/batch.rs @@ -17,7 +17,10 @@ use crate::traits::Queue; /// A handler that processes a batch of jobs atomically. pub trait BatchJobHandler: Send + Sync { - fn handle_batch(&self, jobs: Vec) -> Pin> + Send>>; + fn handle_batch( + &self, + jobs: Vec, + ) -> Pin> + Send>>; } impl BatchJobHandler for F @@ -25,7 +28,10 @@ where F: Fn(Vec) -> Fut + Send + Sync, Fut: std::future::Future> + Send + 'static, { - fn handle_batch(&self, jobs: Vec) -> Pin> + Send>> { + fn handle_batch( + &self, + jobs: Vec, + ) -> Pin> + Send>> { Box::pin((self)(jobs)) } } @@ -85,7 +91,8 @@ impl BatchProcessor { let handler: Arc = Arc::new(handler); let topic_owned = topic.to_string(); - let (tx, mut rx): (mpsc::Sender, mpsc::Receiver) = mpsc::channel(config.batch_size); + let (tx, mut rx): (mpsc::Sender, mpsc::Receiver) = + mpsc::channel(config.batch_size); // Dequeue loop → forward to channel { @@ -147,7 +154,9 @@ impl BatchProcessor { if !batch.is_empty() { Self::flush(&h, &semaphore, batch).await; } - deadline.as_mut().reset(tokio::time::Instant::now() + config.batch_timeout); + deadline + .as_mut() + .reset(tokio::time::Instant::now() + config.batch_timeout); } }); @@ -210,7 +219,10 @@ mod tests { }); for i in 0..3 { - bp.queue.enqueue("t", format!("job{}", i).into_bytes()).await.unwrap(); + bp.queue + .enqueue("t", format!("job{}", i).into_bytes()) + .await + .unwrap(); } tokio::time::sleep(Duration::from_millis(300)).await; diff --git a/crates/mytheclipse-queue/src/in_memory.rs b/crates/mytheclipse-queue/src/in_memory.rs index 61ee6e6..49a1295 100644 --- a/crates/mytheclipse-queue/src/in_memory.rs +++ b/crates/mytheclipse-queue/src/in_memory.rs @@ -18,7 +18,10 @@ struct TopicQueue { impl std::fmt::Debug for TopicQueue { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("TopicQueue") - .field("jobs_len", &self.jobs.try_lock().map(|j| j.len()).unwrap_or(0)) + .field( + "jobs_len", + &self.jobs.try_lock().map(|j| j.len()).unwrap_or(0), + ) .finish() } } @@ -80,7 +83,10 @@ impl InMemoryQueue { impl Queue for InMemoryQueue { async fn enqueue(&self, topic: &str, payload: Vec) -> Result<(), QueueError> { let tq = self.get_topic(topic).await; - tq.jobs.lock().await.push(Job::new(JobId::generate(), topic, payload)); + tq.jobs + .lock() + .await + .push(Job::new(JobId::generate(), topic, payload)); tq.notify.notify_one(); Ok(()) } @@ -140,7 +146,11 @@ mod tests { async fn enqueue_dequeue_roundtrip() { let q = InMemoryQueue::new(); q.enqueue("test", b"hello".to_vec()).await.unwrap(); - let job = q.dequeue("test", Duration::from_millis(500)).await.unwrap().unwrap(); + let job = q + .dequeue("test", Duration::from_millis(500)) + .await + .unwrap() + .unwrap(); assert_eq!(job.payload, b"hello"); assert_eq!(job.topic, "test"); } diff --git a/crates/mytheclipse-queue/src/lib.rs b/crates/mytheclipse-queue/src/lib.rs index a10f47c..640de29 100644 --- a/crates/mytheclipse-queue/src/lib.rs +++ b/crates/mytheclipse-queue/src/lib.rs @@ -11,28 +11,31 @@ //! - **PostgreSQL** (`postgres`) — `SKIP LOCKED` polling. pub mod error; -pub mod job; #[cfg(feature = "in-memory")] pub mod in_memory; +pub mod job; pub mod traits; +#[cfg(feature = "in-memory")] pub mod worker; -#[cfg(feature = "in-memory")] -pub mod batch; #[cfg(feature = "in-memory")] pub mod backpressure_enqueue; #[cfg(feature = "in-memory")] +pub mod batch; +#[cfg(feature = "in-memory")] pub mod rate_limited; #[cfg(feature = "in-memory")] pub mod worker_rate_limited; #[cfg(feature = "in-memory")] -pub use backpressure_enqueue::{BackpressureEnforcer, BackpressureError, enqueue_with_backpressure}; +pub use backpressure_enqueue::{ + enqueue_with_backpressure, BackpressureEnforcer, BackpressureError, +}; #[cfg(feature = "in-memory")] -pub use rate_limited::{RateLimitedQueue, RateLimitQueueError}; +pub use rate_limited::{RateLimitQueueError, RateLimitedQueue}; #[cfg(feature = "in-memory")] pub use worker_rate_limited::RateLimitedWorkerPool; #[cfg(feature = "in-memory")] pub mod pipeline; #[cfg(feature = "in-memory")] -pub use pipeline::{StageRunner, Stage, StageError}; +pub use pipeline::{Stage, StageError, StageRunner}; diff --git a/crates/mytheclipse-queue/src/pipeline.rs b/crates/mytheclipse-queue/src/pipeline.rs index 54405a5..44ee85a 100644 --- a/crates/mytheclipse-queue/src/pipeline.rs +++ b/crates/mytheclipse-queue/src/pipeline.rs @@ -75,16 +75,14 @@ where let mut input = input; loop { match input.recv().await { - Some(item) => { - match stage.process(item).await { - Ok(out) => { - if output.send(out).await.is_err() { - return Err(StageError::ChannelClosed); - } + Some(item) => match stage.process(item).await { + Ok(out) => { + if output.send(out).await.is_err() { + return Err(StageError::ChannelClosed); } - Err(e) => return Err(e), } - } + Err(e) => return Err(e), + }, None => return Ok(()), } } diff --git a/crates/mytheclipse-queue/src/rate_limited.rs b/crates/mytheclipse-queue/src/rate_limited.rs index 9c0d0a3..7488e83 100644 --- a/crates/mytheclipse-queue/src/rate_limited.rs +++ b/crates/mytheclipse-queue/src/rate_limited.rs @@ -111,7 +111,11 @@ impl Queue for RateLimitedQueue { self.inner.enqueue(topic, payload).await } - async fn dequeue(&self, topic: &str, timeout: Duration) -> Result, QueueError> { + async fn dequeue( + &self, + topic: &str, + timeout: Duration, + ) -> Result, QueueError> { self.inner.dequeue(topic, timeout).await } @@ -119,7 +123,11 @@ impl Queue for RateLimitedQueue { self.inner.ack(job).await } - async fn nack(&self, job: &crate::job::Job, requeue: bool) -> Result<(), crate::error::JobError> { + async fn nack( + &self, + job: &crate::job::Job, + requeue: bool, + ) -> Result<(), crate::error::JobError> { self.inner.nack(job, requeue).await } @@ -148,7 +156,7 @@ mod tests { async fn rejects_when_bucket_empty() { let inner = InMemoryQueue::new(); let rl = RateLimitedQueue::new(inner, 0.0, 1); // 0 tokens/sec, 1 burst - // consume the single burst token + // consume the single burst token let _ = rl.enqueue("t", b"x".to_vec()).await; // next should be rate limited (no refill) let result = rl.enqueue("t", b"y".to_vec()).await; diff --git a/crates/mytheclipse-queue/src/traits.rs b/crates/mytheclipse-queue/src/traits.rs index 0462ba9..a90d5eb 100644 --- a/crates/mytheclipse-queue/src/traits.rs +++ b/crates/mytheclipse-queue/src/traits.rs @@ -3,8 +3,8 @@ use async_trait::async_trait; use std::time::Duration; +use crate::error::{JobError, QueueError}; use crate::job::Job; -use crate::error::{QueueError, JobError}; /// A handle to a single unit of queued work. /// diff --git a/crates/mytheclipse-queue/src/worker.rs b/crates/mytheclipse-queue/src/worker.rs index 18bc5ef..dcf3812 100644 --- a/crates/mytheclipse-queue/src/worker.rs +++ b/crates/mytheclipse-queue/src/worker.rs @@ -71,10 +71,13 @@ pub struct WorkerPool { impl WorkerPool { /// Creates a new worker pool with the given concurrency. pub fn new(queue: Q, concurrency: usize) -> Self { - Self::with_config(queue, WorkerConfig { - concurrency, - ..Default::default() - }) + Self::with_config( + queue, + WorkerConfig { + concurrency, + ..Default::default() + }, + ) } /// Creates a new worker pool with explicit configuration. @@ -148,8 +151,8 @@ impl WorkerPool { /// Computes the (capped) exponential backoff delay. pub fn retry_delay(config: &WorkerConfig, attempt: u32) -> Duration { let exponent = attempt as f64; - let computed = config.retry_base_delay.as_millis() as f64 - * config.retry_factor.powf(exponent.max(0.0)); + let computed = + config.retry_base_delay.as_millis() as f64 * config.retry_factor.powf(exponent.max(0.0)); let capped = computed.min(config.retry_max_delay.as_millis() as f64); Duration::from_millis(capped as u64) } diff --git a/crates/mytheclipse-queue/src/worker_rate_limited.rs b/crates/mytheclipse-queue/src/worker_rate_limited.rs index b094db3..ea563b9 100644 --- a/crates/mytheclipse-queue/src/worker_rate_limited.rs +++ b/crates/mytheclipse-queue/src/worker_rate_limited.rs @@ -6,8 +6,8 @@ //! service faster than its rate limit allows. use crate::rate_limited::RateLimitedQueue; -use crate::worker::{JobHandler, WorkerConfig, WorkerPool}; use crate::traits::Queue; +use crate::worker::{JobHandler, WorkerConfig, WorkerPool}; /// A `WorkerPool` whose dequeue is rate-limited via a token bucket. pub struct RateLimitedWorkerPool { @@ -40,12 +40,8 @@ mod tests { #[test] fn constructs_rate_limited_pool() { use crate::in_memory::InMemoryQueue; - let _pool = RateLimitedWorkerPool::new( - InMemoryQueue::new(), - WorkerConfig::default(), - 10.0, - 5, - ); + let _pool = + RateLimitedWorkerPool::new(InMemoryQueue::new(), WorkerConfig::default(), 10.0, 5); // smoke: just verifies construction } } diff --git a/crates/mytheclipse-queue/tests/rate_limited_stress.rs b/crates/mytheclipse-queue/tests/rate_limited_stress.rs index 474506d..23b2ee3 100644 --- a/crates/mytheclipse-queue/tests/rate_limited_stress.rs +++ b/crates/mytheclipse-queue/tests/rate_limited_stress.rs @@ -36,7 +36,11 @@ async fn rate_limited_queue_no_item_loss_under_contention() { let s = Arc::clone(&seen); workers.push(tokio::spawn(async move { loop { - match q.dequeue("stress", Duration::from_millis(50)).await.unwrap() { + match q + .dequeue("stress", Duration::from_millis(50)) + .await + .unwrap() + { Some(job) => { let _ = String::from_utf8(job.payload).unwrap(); s.fetch_add(1, Ordering::SeqCst); @@ -64,4 +68,4 @@ async fn rate_limited_queue_no_item_loss_under_contention() { TASKS * PER_TASK, "items lost under contention" ); -} \ No newline at end of file +} diff --git a/crates/mytheclipse-storage/src/local.rs b/crates/mytheclipse-storage/src/local.rs index 29d9b48..1c0bc64 100644 --- a/crates/mytheclipse-storage/src/local.rs +++ b/crates/mytheclipse-storage/src/local.rs @@ -67,12 +67,10 @@ impl StorageDriver for LocalFileStorage { let mut file = tokio::fs::File::create(&tmp) .await .map_err(|e| StorageError::Io(e.to_string()))?; - let written = tokio::io::copy(&mut data, &mut file) - .await - .map_err(|e| { - let _ = std::fs::remove_file(&tmp); - StorageError::Io(e.to_string()) - })?; + let written = tokio::io::copy(&mut data, &mut file).await.map_err(|e| { + let _ = std::fs::remove_file(&tmp); + StorageError::Io(e.to_string()) + })?; // Ensure durability: flush to OS, fsync, then rename. tokio::fs::File::open(&tmp) .await diff --git a/crates/mytheclipse-tracing/src/fmt.rs b/crates/mytheclipse-tracing/src/fmt.rs index 1d3eb80..0e9add7 100644 --- a/crates/mytheclipse-tracing/src/fmt.rs +++ b/crates/mytheclipse-tracing/src/fmt.rs @@ -14,9 +14,7 @@ impl TracingLayer { pub fn install() { let filter = EnvFilter::try_from_default_env() .unwrap_or_else(|_| EnvFilter::new("mytheclipse=info")); - let _ = tracing_subscriber::fmt() - .with_env_filter(filter) - .try_init(); + let _ = tracing_subscriber::fmt().with_env_filter(filter).try_init(); } /// Returns a formatted layer for manual composition. diff --git a/crates/mytheclipse-tracing/src/lib.rs b/crates/mytheclipse-tracing/src/lib.rs index 4fdf2d4..ef1d26f 100644 --- a/crates/mytheclipse-tracing/src/lib.rs +++ b/crates/mytheclipse-tracing/src/lib.rs @@ -3,9 +3,21 @@ //! Pre-built tracing layers combining all mytheclipse primitives with //! optional export backends (OTLP, Jaeger, Zipkin). +#[cfg(any( + feature = "env", + feature = "otel", + feature = "jaeger", + feature = "zipkin" +))] pub mod fmt; pub mod otel; +#[cfg(any( + feature = "env", + feature = "otel", + feature = "jaeger", + feature = "zipkin" +))] pub use fmt::TracingLayer; #[cfg(any(feature = "otel", feature = "jaeger", feature = "full"))] pub use otel::OtelLayer; diff --git a/crates/mytheclipse/Cargo.toml b/crates/mytheclipse/Cargo.toml index 90cdd50..1dfe1bc 100644 --- a/crates/mytheclipse/Cargo.toml +++ b/crates/mytheclipse/Cargo.toml @@ -43,11 +43,21 @@ name = "main" path = "examples/main.rs" required-features = ["full"] +[[example]] +name = "high_level" +path = "examples/high_level.rs" +required-features = ["full"] + [[example]] name = "scaling_demo" path = "examples/scaling_demo.rs" required-features = ["full"] +[[test]] +name = "race_stress" +path = "tests/race_stress.rs" +required-features = ["full"] + [[bench]] name = "primitives" path = "benches/primitives.rs" diff --git a/crates/mytheclipse/benches/primitives.rs b/crates/mytheclipse/benches/primitives.rs index 93a9e8d..16c6257 100644 --- a/crates/mytheclipse/benches/primitives.rs +++ b/crates/mytheclipse/benches/primitives.rs @@ -18,8 +18,8 @@ use mytheclipse::aggregate_error::AggregateError; use mytheclipse::parallel_map::{parallel_for_each, parallel_map}; use mytheclipse::pool::{Pool, SemaphorePool}; use mytheclipse::ratelimit::RateLimiter; -use mytheclipse::retry_ext::RetryExt; use mytheclipse::retry::RetryConfig; +use mytheclipse::retry_ext::RetryExt; use mytheclipse::shutdown_guard::ShutdownGuard; fn rt() -> Runtime { @@ -52,9 +52,13 @@ fn bench_parallel_for_each(c: &mut Criterion) { let rt = rt(); c.bench_function("parallel_for_each/1000x8", |b| { b.to_async(&rt).iter(|| async { - parallel_for_each(0u32..1000, 8, |_| async move { Ok::<_, std::io::Error>(()) }) - .await - .unwrap(); + parallel_for_each( + 0u32..1000, + 8, + |_| async move { Ok::<_, std::io::Error>(()) }, + ) + .await + .unwrap(); black_box(()); }); }); @@ -177,4 +181,4 @@ criterion_group!( bench_aggregate_error, bench_shutdown_guard, ); -criterion_main!(benches); \ No newline at end of file +criterion_main!(benches); diff --git a/crates/mytheclipse/examples/high_level.rs b/crates/mytheclipse/examples/high_level.rs index 0d9a4c8..771ec24 100644 --- a/crates/mytheclipse/examples/high_level.rs +++ b/crates/mytheclipse/examples/high_level.rs @@ -21,8 +21,8 @@ use mytheclipse::{ pool::{AutoReconnectPool, Pool, Reconnectable, SemaphorePool}, retry_ext::RetryExt, runtime_auto::RuntimeConfig, - shutdown_guard::ShutdownGuard, service_builder::ServiceConfig, + shutdown_guard::ShutdownGuard, }; #[tokio::main] @@ -66,7 +66,10 @@ async fn main() { } }; let value = fut.retry(cfg, |_: &String| true, op).await.unwrap(); - println!("3. RetryExt with {} attempts -> {value}", attempts.load(Ordering::SeqCst)); + println!( + "3. RetryExt with {} attempts -> {value}", + attempts.load(Ordering::SeqCst) + ); // 4. RAII ShutdownGuard — callback runs exactly once on drop, panic-safe. let fired = Arc::new(AtomicU32::new(0)); @@ -93,18 +96,20 @@ async fn main() { let svc = AutoMetricsServiceBuilder::new("demo_op", svc_cfg); let n = Arc::new(AtomicU32::new(0)); let n2 = Arc::clone(&n); - let _: Result> = svc.run(|| { - let n2 = Arc::clone(&n2); - Box::pin(async move { - tokio::time::sleep(Duration::from_millis(5)).await; - let v = n2.fetch_add(1, Ordering::SeqCst); - if v < 1 { - Err(()) - } else { - Ok(7u32) - } + let _: Result> = svc + .run(|| { + let n2 = Arc::clone(&n2); + Box::pin(async move { + tokio::time::sleep(Duration::from_millis(5)).await; + let v = n2.fetch_add(1, Ordering::SeqCst); + if v < 1 { + Err(()) + } else { + Ok(7u32) + } + }) }) - }).await; + .await; let snap = svc.collector().snapshot(); println!( "6. AutoMetrics -> {} counters, {} histograms", diff --git a/crates/mytheclipse/examples/scaling_demo.rs b/crates/mytheclipse/examples/scaling_demo.rs index 0706ffc..24906ae 100644 --- a/crates/mytheclipse/examples/scaling_demo.rs +++ b/crates/mytheclipse/examples/scaling_demo.rs @@ -16,7 +16,10 @@ use mytheclipse::parallel_map::{parallel_for_each, ParallelConcurrency}; #[tokio::main] async fn main() { let total = 100u32; - println!("host available_parallelism = {}", <() as ParallelConcurrency>::resolve(())); + println!( + "host available_parallelism = {}", + <() as ParallelConcurrency>::resolve(()) + ); println!("total items = {total}"); println!(); @@ -49,6 +52,12 @@ async fn run(label: &str, concurrency: impl ParallelConcurrency + Copy, total: u let elapsed = start.elapsed(); println!("{label} resolved={resolved}"); - println!(" peak in-flight = {} (bounded, never {total})", peak.load(Ordering::SeqCst)); - println!(" elapsed = {elapsed:?} (sequential ~{}ms)", total * 1); -} \ No newline at end of file + println!( + " peak in-flight = {} (bounded, never {total})", + peak.load(Ordering::SeqCst) + ); + println!( + " elapsed = {elapsed:?} (sequential ~{}ms)", + total * 1 + ); +} diff --git a/crates/mytheclipse/src/aggregate_error.rs b/crates/mytheclipse/src/aggregate_error.rs index 86f0061..dcddd7c 100644 --- a/crates/mytheclipse/src/aggregate_error.rs +++ b/crates/mytheclipse/src/aggregate_error.rs @@ -117,12 +117,18 @@ impl std::error::Error for AggregateError {} impl From>> for AggregateError { fn from(errors: Vec>) -> Self { - Self { errors, context: None } + Self { + errors, + context: None, + } } } impl Extend> for AggregateError { - fn extend>>(&mut self, iter: T) { + fn extend>>( + &mut self, + iter: T, + ) { self.errors.extend(iter); } } @@ -160,8 +166,7 @@ mod tests { #[test] fn collects_values_when_all_ok() { - let results: Vec> = - vec![Ok(1), Ok(2), Ok(3)]; + let results: Vec> = vec![Ok(1), Ok(2), Ok(3)]; let out = AggregateError::from_results(results).unwrap(); assert_eq!(out, vec![1, 2, 3]); } diff --git a/crates/mytheclipse/src/auto_metrics_service.rs b/crates/mytheclipse/src/auto_metrics_service.rs index c5f0ab8..a3c1f40 100644 --- a/crates/mytheclipse/src/auto_metrics_service.rs +++ b/crates/mytheclipse/src/auto_metrics_service.rs @@ -74,8 +74,8 @@ impl AutoMetricsServiceBuilder { self } - /// Attaches a [`MetricsBridge`] to forward snapshots downstream (requires - /// the `resiliency` feature which pulls in the bridge). + /// Attaches a [`crate::metrics_bridge::MetricsBridge`] to forward snapshots + /// downstream (requires the `resiliency` feature which pulls in the bridge). #[cfg(feature = "resiliency")] pub fn with_bridge(mut self, bridge: crate::metrics_bridge::MetricsBridge) -> Self { self.bridge = Some(bridge); @@ -108,8 +108,13 @@ impl AutoMetricsServiceBuilder { Err(_) => "other", }; - self.metrics - .inc_counter(&format!("mytheclipse_service_calls_total{{service=\"{}\",outcome=\"{}\"}}", self.service_name, outcome), 1); + self.metrics.inc_counter( + &format!( + "mytheclipse_service_calls_total{{service=\"{}\",outcome=\"{}\"}}", + self.service_name, outcome + ), + 1, + ); self.metrics .observe("mytheclipse_service_duration_seconds", dur); @@ -125,8 +130,8 @@ impl AutoMetricsServiceBuilder { #[cfg(test)] mod tests { use super::*; + use std::sync::atomic::{AtomicU32, Ordering}; use std::sync::Arc; - use std::sync::atomic::{Ordering, AtomicU32}; #[tokio::test] async fn auto_metrics_records_call() { @@ -136,13 +141,19 @@ mod tests { let attempts = Arc::new(AtomicU32::new(0)); let a = Arc::clone(&attempts); - let result: Result> = builder.run(|| { - let a = Arc::clone(&a); - Box::pin(async move { - let n = a.fetch_add(1, Ordering::SeqCst); - if n < 2 { Err(()) } else { Ok(42u32) } + let result: Result> = builder + .run(|| { + let a = Arc::clone(&a); + Box::pin(async move { + let n = a.fetch_add(1, Ordering::SeqCst); + if n < 2 { + Err(()) + } else { + Ok(42u32) + } + }) }) - }).await; + .await; assert_eq!(result.unwrap(), 42); assert_eq!(attempts.load(Ordering::SeqCst), 3); let snap = builder.collector().snapshot(); diff --git a/crates/mytheclipse/src/bg_join.rs b/crates/mytheclipse/src/bg_join.rs index db76702..a205348 100644 --- a/crates/mytheclipse/src/bg_join.rs +++ b/crates/mytheclipse/src/bg_join.rs @@ -8,7 +8,6 @@ //! *concurrency*, `BgJoiner` adds structured *lifetimes* so a service can wait //! for all in-flight work to settle before terminating. -use std::future::Future; use std::sync::Arc; use std::time::Duration; @@ -36,7 +35,9 @@ impl BgJoiner { F: std::future::Future + Send + 'static, F::Output: Send + 'static, { - let handle: JoinHandle<()> = tokio::spawn(async move { let _ = future.await; }); + let handle: JoinHandle<()> = tokio::spawn(async move { + let _ = future.await; + }); self.track(handle); } @@ -55,6 +56,11 @@ impl BgJoiner { self.inner.lock().await.len() } + /// Returns `true` if there are no currently-tracked tasks. + pub async fn is_empty(&self) -> bool { + self.inner.lock().await.is_empty() + } + /// Await every tracked task, dropping any that are still pending once /// `deadline` elapses. Returns the count of tasks that had not completed /// within the timeout. @@ -66,7 +72,6 @@ impl BgJoiner { }; let mut pending: Vec> = handles; - let mut dropped = 0usize; loop { if pending.is_empty() { @@ -74,11 +79,11 @@ impl BgJoiner { } if now.elapsed() >= deadline { - dropped = pending.len(); + let count = pending.len(); for h in pending.drain(..) { h.abort(); } - return dropped; + return count; } let remaining = deadline.saturating_sub(now.elapsed()); diff --git a/crates/mytheclipse/src/compute.rs b/crates/mytheclipse/src/compute.rs index 5a96f2a..5527fb8 100644 --- a/crates/mytheclipse/src/compute.rs +++ b/crates/mytheclipse/src/compute.rs @@ -121,18 +121,16 @@ where F: Fn(I::Item) -> Result + Send + Sync, { let wrapped = AssertUnwindSafe(f); - let collected: Vec> = context() - .compute_pool - .install(move || { - let f = wrapped; - items - .into_par_iter() - .map(|item| { - catch_unwind(AssertUnwindSafe(|| f(item))) - .unwrap_or_else(|payload| Err(panic_payload_to_string(payload))) - }) - .collect() - }); + let collected: Vec> = context().compute_pool.install(move || { + let f = wrapped; + items + .into_par_iter() + .map(|item| { + catch_unwind(AssertUnwindSafe(|| f(item))) + .unwrap_or_else(|payload| Err(panic_payload_to_string(payload))) + }) + .collect() + }); let mut values = Vec::with_capacity(collected.len()); let mut errors = Vec::new(); @@ -165,10 +163,7 @@ where /// ).unwrap(); /// assert_eq!(a + b, (0..2_000_000u64).sum::()); /// ``` -pub fn compute_join( - a: A, - b: B, -) -> Result<(RA, RB), MytheclipseError> +pub fn compute_join(a: A, b: B) -> Result<(RA, RB), MytheclipseError> where A: FnOnce() -> RA + Send, RA: Send, @@ -177,15 +172,19 @@ where { let a = AssertUnwindSafe(a); let b = AssertUnwindSafe(b); - context() - .compute_pool - .install(|| { - let (ra, rb) = rayon::join( - move || catch_unwind(a).map_err(|p| MytheclipseError::ComputePanic(panic_payload_to_string(p))), - move || catch_unwind(b).map_err(|p| MytheclipseError::ComputePanic(panic_payload_to_string(p))), - ); - Ok((ra?, rb?)) - }) + context().compute_pool.install(|| { + let (ra, rb) = rayon::join( + move || { + catch_unwind(a) + .map_err(|p| MytheclipseError::ComputePanic(panic_payload_to_string(p))) + }, + move || { + catch_unwind(b) + .map_err(|p| MytheclipseError::ComputePanic(panic_payload_to_string(p))) + }, + ); + Ok((ra?, rb?)) + }) } /// Runs `f` over every item on the compute pool in parallel, discarding @@ -211,25 +210,20 @@ where F: Fn(I::Item) -> Result<(), String> + Send + Sync, { let wrapped = AssertUnwindSafe(f); - let collected: Vec> = context() - .compute_pool - .install(move || { - let f = wrapped; - items - .into_par_iter() - .map(|item| { - catch_unwind(AssertUnwindSafe(|| f(item))) - .unwrap_or_else(|payload| Err(panic_payload_to_string(payload))) - }) - .collect() - }); + let collected: Vec> = context().compute_pool.install(move || { + let f = wrapped; + items + .into_par_iter() + .map(|item| { + catch_unwind(AssertUnwindSafe(|| f(item))) + .unwrap_or_else(|payload| Err(panic_payload_to_string(payload))) + }) + .collect() + }); if collected.iter().any(|r| r.is_err()) { Err(ComputeErrors { - errors: collected - .into_iter() - .filter_map(|r| r.err()) - .collect(), + errors: collected.into_iter().filter_map(|r| r.err()).collect(), }) } else { Ok(()) diff --git a/crates/mytheclipse/src/dlock.rs b/crates/mytheclipse/src/dlock.rs index 0f31b54..7c60430 100644 --- a/crates/mytheclipse/src/dlock.rs +++ b/crates/mytheclipse/src/dlock.rs @@ -69,7 +69,12 @@ impl Drop for LockGuard { #[async_trait] pub trait DistributedLock: Send + Sync { /// Attempts to acquire the lock with the given lease duration. - async fn acquire(&self, key: &str, lease: Duration, timeout: Duration) -> Result; + async fn acquire( + &self, + key: &str, + lease: Duration, + timeout: Duration, + ) -> Result; /// Releases the lock. async fn release(&self, key: &str) -> Result<(), LockError>; @@ -90,14 +95,6 @@ impl InProcLock { held: Arc::new(Mutex::new(std::collections::HashMap::new())), } } - - fn is_expired(map: &std::collections::HashMap, key: &str) -> bool { - if let Some(expiry) = map.get(key) { - *expiry <= Instant::now() - } else { - false - } - } } impl Default for InProcLock { @@ -108,7 +105,12 @@ impl Default for InProcLock { #[async_trait] impl DistributedLock for InProcLock { - async fn acquire(&self, key: &str, lease: Duration, timeout: Duration) -> Result { + async fn acquire( + &self, + key: &str, + lease: Duration, + timeout: Duration, + ) -> Result { let deadline = Instant::now() + timeout; loop { { @@ -154,7 +156,10 @@ mod tests { #[tokio::test] async fn lock_acquire_release() { let lock = InProcLock::new(); - let guard = lock.acquire("key", Duration::from_secs(10), Duration::from_secs(1)).await.unwrap(); + let guard = lock + .acquire("key", Duration::from_secs(10), Duration::from_secs(1)) + .await + .unwrap(); assert!(lock.release("key").await.is_ok()); drop(guard); } @@ -162,9 +167,14 @@ mod tests { #[tokio::test] async fn lock_rejects_second_acquire() { let lock = InProcLock::new(); - let _guard1 = lock.acquire("key", Duration::from_secs(10), Duration::from_secs(1)).await.unwrap(); + let _guard1 = lock + .acquire("key", Duration::from_secs(10), Duration::from_secs(1)) + .await + .unwrap(); // While guard1 is alive, a second acquire with short timeout should fail. - let result = lock.acquire("key", Duration::from_secs(10), Duration::from_millis(50)).await; + let result = lock + .acquire("key", Duration::from_secs(10), Duration::from_millis(50)) + .await; assert!(result.is_err()); drop(_guard1); } @@ -172,21 +182,31 @@ mod tests { #[tokio::test] async fn lock_auto_releases_on_drop() { let lock = InProcLock::new(); - let guard = lock.acquire("k", Duration::from_secs(10), Duration::from_secs(1)).await.unwrap(); + let guard = lock + .acquire("k", Duration::from_secs(10), Duration::from_secs(1)) + .await + .unwrap(); drop(guard); // After drop, the lock should be releasable / re-acquirable. - let result = lock.acquire("k", Duration::from_secs(10), Duration::from_millis(50)).await; + let result = lock + .acquire("k", Duration::from_secs(10), Duration::from_millis(50)) + .await; assert!(result.is_ok(), "lock should be free after guard drop"); } #[tokio::test] async fn lock_expires_after_lease() { let lock = InProcLock::new(); - let _guard = lock.acquire("key", Duration::from_millis(20), Duration::from_millis(5)).await.unwrap(); + let _guard = lock + .acquire("key", Duration::from_millis(20), Duration::from_millis(5)) + .await + .unwrap(); drop(_guard); tokio::time::sleep(Duration::from_millis(30)).await; // Should be acquirable now. - let result = lock.acquire("key", Duration::from_millis(20), Duration::from_millis(5)).await; + let result = lock + .acquire("key", Duration::from_millis(20), Duration::from_millis(5)) + .await; assert!(result.is_ok()); } } diff --git a/crates/mytheclipse/src/health.rs b/crates/mytheclipse/src/health.rs index 7e9ad5f..1b5a5a5 100644 --- a/crates/mytheclipse/src/health.rs +++ b/crates/mytheclipse/src/health.rs @@ -26,7 +26,9 @@ impl fmt::Display for HealthStatus { /// A single health check. pub trait HealthCheck: Send + Sync { fn name(&self) -> &str; - fn check(&self) -> std::pin::Pin + Send + '_>>; + fn check( + &self, + ) -> std::pin::Pin + Send + '_>>; } /// A registered health check with its name and trait object. diff --git a/crates/mytheclipse/src/lib.rs b/crates/mytheclipse/src/lib.rs index b537477..ae22403 100644 --- a/crates/mytheclipse/src/lib.rs +++ b/crates/mytheclipse/src/lib.rs @@ -12,7 +12,7 @@ //! - [`spawn_io`] (feature `io`) — spawn an async I/O task, tracing-instrumented. //! - [`compute()`] (feature `compute`) — run CPU-bound work on a sized Rayon pool, panic-isolated. //! - [`spawn_bg`] (feature `bg`) — spawn a background task under bounded concurrency. -//! - [`retry`] / [`CircuitBreaker`] / [`timeout()`] (feature `resiliency`) — fault tolerance. +//! - [`retry()`] / [`CircuitBreaker`] / [`timeout()`] (feature `resiliency`) — fault tolerance. //! - [`RateLimiter`] / [`BackpressureQueue`] / [`ConcurrencyLimiter`] (feature `traffic`) — load control. //! - [`SemaphorePool`] (feature `traffic`) — shared bounded resource pool. //! - [`ShutdownManager`] / [`CronSchedule`] (feature `lifecycle`) — lifecycle + scheduling. @@ -24,57 +24,59 @@ pub mod context; pub mod error; -#[cfg(feature = "io")] -pub mod io; -#[cfg(feature = "compute")] -pub mod compute; #[cfg(feature = "bg")] pub mod bg; +#[cfg(feature = "compute")] +pub mod compute; +#[cfg(feature = "io")] +pub mod io; -#[cfg(feature = "resiliency")] -pub mod retry; -#[cfg(feature = "resiliency")] -pub mod retry_ext; #[cfg(feature = "resiliency")] pub mod aggregate_error; #[cfg(feature = "resiliency")] pub mod parallel_map; #[cfg(feature = "resiliency")] -pub use retry_ext::RetryExt; +pub mod retry; +#[cfg(feature = "resiliency")] +pub mod retry_ext; #[cfg(feature = "resiliency")] pub use aggregate_error::AggregateError; #[cfg(feature = "resiliency")] -pub use parallel_map::{parallel_map, parallel_map_unordered, parallel_for_each, ParallelConcurrency}; -#[cfg(feature = "observability")] +pub use parallel_map::{ + parallel_for_each, parallel_map, parallel_map_unordered, ParallelConcurrency, +}; +#[cfg(feature = "resiliency")] +pub use retry_ext::RetryExt; +#[cfg(all(feature = "observability", feature = "resiliency"))] pub mod auto_metrics_service; -#[cfg(feature = "observability")] +#[cfg(all(feature = "observability", feature = "resiliency"))] pub use auto_metrics_service::AutoMetricsServiceBuilder; #[cfg(feature = "resiliency")] pub mod circuit_breaker; #[cfg(feature = "resiliency")] pub mod timeout; -#[cfg(feature = "traffic")] -pub mod ratelimit; #[cfg(feature = "traffic")] pub mod backpressure; #[cfg(feature = "traffic")] pub mod concurrency; #[cfg(feature = "traffic")] pub mod pool; +#[cfg(feature = "traffic")] +pub mod ratelimit; -#[cfg(feature = "lifecycle")] -pub mod shutdown; #[cfg(feature = "lifecycle")] 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")] pub mod lifecycle; +#[cfg(all(feature = "observability", feature = "traffic"))] +pub mod pool_health; +#[cfg(feature = "lifecycle")] +pub mod shutdown; #[cfg(feature = "lifecycle")] pub mod bg_join; @@ -94,59 +96,59 @@ pub mod middleware; #[cfg(feature = "observability")] pub mod metrics; #[cfg(feature = "observability")] -pub mod panic_tracker; -#[cfg(feature = "observability")] pub mod metrics_bridge; +#[cfg(feature = "observability")] +pub mod panic_tracker; -#[cfg(feature = "resiliency")] -pub mod service_builder; #[cfg(feature = "lifecycle")] pub mod dlock; +#[cfg(feature = "resiliency")] +pub mod service_builder; pub use context::{context, EngineContext}; pub use error::MytheclipseError; -#[cfg(feature = "io")] -pub use io::spawn_io; -#[cfg(feature = "compute")] -pub use compute::{compute, compute_join, compute_map, compute_par_for_each, ComputeErrors}; #[cfg(feature = "bg")] pub use bg::spawn_bg; +#[cfg(feature = "compute")] +pub use compute::{compute, compute_join, compute_map, compute_par_for_each, ComputeErrors}; +#[cfg(feature = "io")] +pub use io::spawn_io; -#[cfg(feature = "resiliency")] -pub use retry::{retry, JitterKind, RetryConfig, RetryError}; #[cfg(feature = "resiliency")] pub use circuit_breaker::{CircuitBreaker, CircuitBreakerConfig, CircuitError, CircuitState}; #[cfg(feature = "resiliency")] +pub use retry::{retry, JitterKind, RetryConfig, RetryError}; +#[cfg(feature = "resiliency")] pub use timeout::{timeout, with_timeout, Timeout, TimeoutError}; -#[cfg(feature = "traffic")] -pub use ratelimit::{RateLimitError, RateLimiter}; +/// Re-export of `async-trait` so implementing [`pool::Reconnectable`] (and +/// other async traits) doesn't require users to add their own `async-trait` +/// dependency. +pub use async_trait::async_trait; #[cfg(feature = "traffic")] pub use backpressure::{BackpressureError, BackpressureQueue, OverflowPolicy}; #[cfg(feature = "traffic")] pub use concurrency::{ConcurrencyLimiter, ConcurrencyPermit}; #[cfg(feature = "traffic")] -pub use pool::{Pool, PoolError, Pooled, SemaphorePool, AutoReconnectPool, Reconnectable}; -/// Re-export of `async-trait` so implementing [`pool::Reconnectable`] (and -/// other async traits) doesn't require users to add their own `async-trait` -/// dependency. -pub use async_trait::async_trait; +pub use pool::{AutoReconnectPool, Pool, PoolError, Pooled, Reconnectable, SemaphorePool}; +#[cfg(feature = "traffic")] +pub use ratelimit::{RateLimitError, RateLimiter}; -#[cfg(feature = "lifecycle")] -pub use shutdown::{ShutdownManager, ShutdownSignal}; #[cfg(feature = "lifecycle")] pub use cron::{schedule, CronError, CronJob, CronParseError, CronSchedule}; #[cfg(feature = "lifecycle")] pub use health::{HealthCheck, HealthRegistry, HealthStatus}; #[cfg(feature = "lifecycle")] pub use leader::{InProcLeaderElection, LeaderElection}; +#[cfg(feature = "lifecycle")] +pub use shutdown::{ShutdownManager, ShutdownSignal}; #[cfg(feature = "resiliency")] pub use service_builder::{RunError, ServiceBuilder, ServiceConfig}; #[cfg(feature = "lifecycle")] -pub use dlock::{DistributedLock, LockError, LockGuard, InProcLock}; +pub use dlock::{DistributedLock, InProcLock, LockError, LockGuard}; #[cfg(feature = "lifecycle")] pub use lifecycle::AsyncLifecycleManager; @@ -169,7 +171,7 @@ pub use pool_health::HealthCheckedPool; pub use bg_join::BgJoiner; #[cfg(all(feature = "observability", feature = "resiliency"))] -pub use middleware::{MiddlewarePipeline, PipelineError, BoxMiddleware, mw}; +pub use middleware::{mw, BoxMiddleware, MiddlewarePipeline, PipelineError}; #[cfg(feature = "observability")] pub use panic_tracker::{PanicGuard, PanicInfo, PanicTracker}; diff --git a/crates/mytheclipse/src/lifecycle.rs b/crates/mytheclipse/src/lifecycle.rs index 48183cf..bb42a6c 100644 --- a/crates/mytheclipse/src/lifecycle.rs +++ b/crates/mytheclipse/src/lifecycle.rs @@ -15,9 +15,7 @@ use crate::shutdown::ShutdownManager; /// /// Typical usage: /// ```ignore -/// # tokio::runtime::Runtime::new().unwrap().block_on(async { -/// # use mytheclipse::AsyncLifecycleManager; -/// let mgr = AsyncLifecycleManager::new(); +/// let mgr = mytheclipse::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; @@ -47,7 +45,11 @@ impl AsyncLifecycleManager { } /// Registers a named health check. - pub async fn register_health_check(&self, name: impl Into, check: impl HealthCheck + 'static) { + pub async fn register_health_check( + &self, + name: impl Into, + check: impl HealthCheck + 'static, + ) { self.health.register(name, check).await; } @@ -73,7 +75,7 @@ impl AsyncLifecycleManager { loop { // Stop when shutdown is requested. if sig.is_shutdown() { - tracing::info_span!("mytheclipse_health_loop", ); + tracing::info_span!("mytheclipse_health_loop",); return; } tokio::select! { @@ -124,16 +126,26 @@ mod tests { struct AlwaysOk; impl HealthCheck for AlwaysOk { - fn name(&self) -> &str { "always-ok" } - fn check(&self) -> std::pin::Pin + Send + '_>> { + 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 + '_>> { + fn name(&self) -> &str { + "always-bad" + } + fn check( + &self, + ) -> std::pin::Pin + Send + '_>> + { Box::pin(async { HealthStatus::Unhealthy }) } } diff --git a/crates/mytheclipse/src/metrics_bridge.rs b/crates/mytheclipse/src/metrics_bridge.rs index a76b8e7..e7799df 100644 --- a/crates/mytheclipse/src/metrics_bridge.rs +++ b/crates/mytheclipse/src/metrics_bridge.rs @@ -10,8 +10,8 @@ use std::time::Duration; use crate::health::{HealthCheck, HealthStatus}; use crate::metrics::MetricsCollector; -/// A health check backed by a [`CircuitBreaker`]: unhealthy if open, -/// degraded if half-open, ok otherwise. +/// A health check backed by a [`crate::circuit_breaker::CircuitBreaker`]: +/// unhealthy if open, degraded if half-open, ok otherwise. /// /// Only available when both `resiliency` and `observability` features are /// enabled (circuit breaker + health/metrics bridge). @@ -33,7 +33,9 @@ impl HealthCheck for CircuitBreakerHealthCheck { "circuit_breaker" } - fn check(&self) -> std::pin::Pin + Send + '_>> { + fn check( + &self, + ) -> std::pin::Pin + Send + '_>> { let state = self.breaker.snapshot().state; Box::pin(async move { match state { @@ -63,11 +65,7 @@ impl MetricsHealthCheck { } fn has_errors(&self) -> bool { - self.collector - .snapshot() - .counters - .values() - .any(|&v| v > 0) + self.collector.snapshot().counters.values().any(|&v| v > 0) } } @@ -76,7 +74,9 @@ impl HealthCheck for MetricsHealthCheck { "metrics" } - fn check(&self) -> std::pin::Pin + Send + '_>> { + fn check( + &self, + ) -> std::pin::Pin + Send + '_>> { let has_errors = self.has_errors(); Box::pin(async move { if has_errors { @@ -177,7 +177,7 @@ mod tests { bridge.emit_now(); } -#[cfg(feature = "lifecycle")] + #[cfg(feature = "lifecycle")] #[tokio::test] async fn lifecycle_manager_with_metrics_bridge() { let collector = MetricsCollector::new(); diff --git a/crates/mytheclipse/src/middleware.rs b/crates/mytheclipse/src/middleware.rs index a203f89..5b1f317 100644 --- a/crates/mytheclipse/src/middleware.rs +++ b/crates/mytheclipse/src/middleware.rs @@ -24,11 +24,8 @@ impl std::fmt::Display for PipelineError { impl std::error::Error for PipelineError {} /// A single async middleware stage. -pub type BoxMiddleware = Arc< - dyn Fn(S) -> Pin> + Send>> - + Send - + Sync, ->; +pub type BoxMiddleware = + Arc Pin> + Send>> + Send + Sync>; /// Helper to box any `async fn` middleware. pub fn mw(f: F) -> BoxMiddleware @@ -48,7 +45,9 @@ pub struct MiddlewarePipeline { impl MiddlewarePipeline { pub fn new() -> Self { - Self { layers: Arc::new(Mutex::new(Vec::new())) } + Self { + layers: Arc::new(Mutex::new(Vec::new())), + } } /// Appends a middleware stage. @@ -58,7 +57,10 @@ impl MiddlewarePipeline { /// Applies every layer in order, short-circuiting on the first error. pub async fn apply(&self, state: S) -> Result { - let layers = self.layers.lock().unwrap(); + // Clone the Arc'd layers out so the MutexGuard is dropped before any + // await point (holding it across `layer(current).await` is unsound — + // a re-entrant layer could deadlock). + let layers: Vec> = self.layers.lock().unwrap().clone(); let mut current = state; for layer in layers.iter() { current = layer(current).await?; @@ -100,7 +102,9 @@ mod tests { async fn short_circuits_on_error() { let p: MiddlewarePipeline = MiddlewarePipeline::new(); let reject = mw(|_s: String| async { - Err::<_, PipelineError>(PipelineError { msg: "rejected".into() }) + Err::<_, PipelineError>(PipelineError { + msg: "rejected".into(), + }) }); p.add(reject); assert!(matches!(p.apply("x".to_string()).await, Err(_))); diff --git a/crates/mytheclipse/src/parallel_map.rs b/crates/mytheclipse/src/parallel_map.rs index 3b87220..3e75c89 100644 --- a/crates/mytheclipse/src/parallel_map.rs +++ b/crates/mytheclipse/src/parallel_map.rs @@ -234,8 +234,7 @@ where // Producer: feed items into the bounded channel (backpressures when all // workers are busy — no full materialization). tokio::spawn(async move { - let mut it = items.into_iter(); - while let Some(item) = it.next() { + for item in items { if tx.send(item).await.is_err() { break; // all workers dropped } @@ -286,11 +285,9 @@ mod tests { #[tokio::test] async fn maps_in_order_with_bounded_concurrency() { - let out = parallel_map( - vec![1, 2, 3, 4], - 2, - |x: i32| async move { Ok::<_, std::io::Error>(x * 2) }, - ) + let out = parallel_map(vec![1, 2, 3, 4], 2, |x: i32| async move { + Ok::<_, std::io::Error>(x * 2) + }) .await .unwrap(); assert_eq!(out, vec![2, 4, 6, 8]); @@ -298,17 +295,13 @@ mod tests { #[tokio::test] async fn aggregates_errors_from_failing_tasks() { - let out = parallel_map( - vec![1, 2, 3], - 4, - |x: i32| async move { - if x == 2 { - Err(std::io::Error::new(std::io::ErrorKind::Other, "boom")) - } else { - Ok::<_, std::io::Error>(x) - } - }, - ) + let out = parallel_map(vec![1, 2, 3], 4, |x: i32| async move { + if x == 2 { + Err(std::io::Error::new(std::io::ErrorKind::Other, "boom")) + } else { + Ok::<_, std::io::Error>(x) + } + }) .await; assert!(out.is_err()); assert_eq!(out.unwrap_err().len(), 1); @@ -316,12 +309,11 @@ mod tests { #[tokio::test] async fn empty_input_returns_empty() { - let out: Result, AggregateError> = parallel_map( - Vec::::new(), - 4, - |x: i32| async move { Ok::<_, std::io::Error>(x) }, - ) - .await; + let out: Result, AggregateError> = + parallel_map(Vec::::new(), 4, |x: i32| async move { + Ok::<_, std::io::Error>(x) + }) + .await; assert_eq!(out.unwrap(), vec![]); } @@ -330,17 +322,13 @@ mod tests { use std::sync::atomic::{AtomicUsize, Ordering}; let count = Arc::new(AtomicUsize::new(0)); let c = Arc::clone(&count); - let out = parallel_for_each( - vec![1_i32, 2, 3, 4, 5], - 2, - move |_: i32| { - let c = Arc::clone(&c); - async move { - c.fetch_add(1, Ordering::SeqCst); - Ok::<_, std::io::Error>(()) - } - }, - ) + let out = parallel_for_each(vec![1_i32, 2, 3, 4, 5], 2, move |_: i32| { + let c = Arc::clone(&c); + async move { + c.fetch_add(1, Ordering::SeqCst); + Ok::<_, std::io::Error>(()) + } + }) .await; assert!(out.is_ok()); assert_eq!(count.load(Ordering::SeqCst), 5); diff --git a/crates/mytheclipse/src/pool.rs b/crates/mytheclipse/src/pool.rs index 804f2b0..314a022 100644 --- a/crates/mytheclipse/src/pool.rs +++ b/crates/mytheclipse/src/pool.rs @@ -60,7 +60,11 @@ impl SemaphorePool { #[async_trait] impl Pool for SemaphorePool { async fn acquire(&self) -> Result, PoolError> { - let permit = self.semaphore.clone().acquire_owned().await + let permit = self + .semaphore + .clone() + .acquire_owned() + .await .map_err(|_| PoolError::Exhausted)?; let idx = ACQUIRE_COUNT.fetch_add(1, Ordering::Relaxed) % self.items.len(); Ok(Pooled { @@ -158,11 +162,7 @@ where _permit: pooled._permit, }) } else { - let fresh = self - .reconnect - .reconnect() - .await - .map_err(PoolError::Other)?; + let fresh = self.reconnect.reconnect().await.map_err(PoolError::Other)?; Ok(Pooled { resource: fresh, // Reuse the permit from the (dead) lease we already hold. diff --git a/crates/mytheclipse/src/pool_health.rs b/crates/mytheclipse/src/pool_health.rs index 492ea8e..1a51f1a 100644 --- a/crates/mytheclipse/src/pool_health.rs +++ b/crates/mytheclipse/src/pool_health.rs @@ -8,7 +8,7 @@ use std::sync::Arc; use std::time::Duration; use crate::health::{HealthCheck, HealthRegistry, HealthStatus}; -use crate::pool::{Pool, Pooled, PoolError, SemaphorePool}; +use crate::pool::{Pool, PoolError, Pooled, SemaphorePool}; /// A health check backed by a closure. struct ClosureCheck { @@ -21,7 +21,9 @@ impl HealthCheck for ClosureCheck { &self.name } - fn check(&self) -> std::pin::Pin + Send + '_>> { + fn check( + &self, + ) -> std::pin::Pin + Send + '_>> { let status = (self.check)(); Box::pin(async move { status }) } diff --git a/crates/mytheclipse/src/retry.rs b/crates/mytheclipse/src/retry.rs index f8f821d..020e866 100644 --- a/crates/mytheclipse/src/retry.rs +++ b/crates/mytheclipse/src/retry.rs @@ -320,10 +320,19 @@ mod tests { ..RetryConfig::default() }; let calls = Cell::new(0u32); - let (result, stats) = retry_with_stats(config, |_| true, || async { - calls.set(calls.get() + 1); - if calls.get() < 3 { Err::("fail") } else { Ok(42u32) } - }).await; + let (result, stats) = retry_with_stats( + config, + |_| true, + || async { + calls.set(calls.get() + 1); + if calls.get() < 3 { + Err::("fail") + } else { + Ok(42u32) + } + }, + ) + .await; assert_eq!(result.unwrap(), 42); assert_eq!(stats.attempts, 3); assert_eq!(stats.retries, 2); diff --git a/crates/mytheclipse/src/retry_ext.rs b/crates/mytheclipse/src/retry_ext.rs index 8c173da..4f64a43 100644 --- a/crates/mytheclipse/src/retry_ext.rs +++ b/crates/mytheclipse/src/retry_ext.rs @@ -76,21 +76,29 @@ where #[cfg(test)] mod tests { use super::*; - use std::time::Duration; - use std::sync::Arc; use std::sync::atomic::{AtomicU32, Ordering}; + use std::sync::Arc; + use std::time::Duration; #[tokio::test] async fn retry_ext_retries_then_succeeds() { let attempts = Arc::new(AtomicU32::new(0)); let a = Arc::clone(&attempts); - let cfg = RetryConfig { max_attempts: 3, base_delay: Duration::from_millis(1), ..RetryConfig::default() }; + let cfg = RetryConfig { + max_attempts: 3, + base_delay: Duration::from_millis(1), + ..RetryConfig::default() + }; let op = move || { let a = Arc::clone(&a); async move { let n = a.fetch_add(1, Ordering::SeqCst); - if n < 2 { Err::<(), String>("transient".into()) } else { Ok(()) } + if n < 2 { + Err::<(), String>("transient".into()) + } else { + Ok(()) + } } }; diff --git a/crates/mytheclipse/src/service_builder.rs b/crates/mytheclipse/src/service_builder.rs index 3797e56..7ff0899 100644 --- a/crates/mytheclipse/src/service_builder.rs +++ b/crates/mytheclipse/src/service_builder.rs @@ -11,10 +11,10 @@ use tracing::Instrument; #[cfg(feature = "resiliency")] use crate::circuit_breaker::CircuitBreaker; -#[cfg(feature = "resiliency")] -use crate::retry::{retry, RetryConfig, RetryError}; #[cfg(feature = "traffic")] use crate::ratelimit::RateLimiter; +#[cfg(feature = "resiliency")] +use crate::retry::{retry, RetryConfig, RetryError}; /// Error returned by [`ServiceBuilder::run`]. #[derive(Debug)] @@ -55,7 +55,10 @@ pub struct ServiceConfig { #[cfg(not(feature = "traffic"))] impl Default for ServiceConfig { fn default() -> Self { - Self { max_attempts: 0, timeout: Duration::ZERO } + Self { + max_attempts: 0, + timeout: Duration::ZERO, + } } } @@ -71,7 +74,12 @@ pub struct ServiceConfig { #[cfg(feature = "traffic")] impl Default for ServiceConfig { fn default() -> Self { - Self { max_attempts: 0, timeout: Duration::ZERO, rate_per_sec: 0.0, rate_burst: 0 } + Self { + max_attempts: 0, + timeout: Duration::ZERO, + rate_per_sec: 0.0, + rate_burst: 0, + } } } @@ -146,7 +154,11 @@ impl ServiceBuilder { #[cfg(feature = "resiliency")] fn record(&self, ok: bool) { if let Some(cb) = &self.circuit { - if ok { cb.record_success(); } else { cb.record_failure(); } + if ok { + cb.record_success(); + } else { + cb.record_failure(); + } } } @@ -162,7 +174,7 @@ impl ServiceBuilder { #[cfg(feature = "resiliency")] { if let Some(retry_cfg) = &self.retry_cfg { - let mut op = f; + let op = f; let result: Result> = if dur > Duration::ZERO { // We can't easily combine retry + timeout with FnMut due to // closure capture rules, so use a manual retry loop instead: @@ -171,7 +183,8 @@ impl ServiceBuilder { let mut op_ref = op; loop { attempt_no += 1; - let span = tracing::info_span!("mytheclipse_service_call", attempt = attempt_no); + let span = + tracing::info_span!("mytheclipse_service_call", attempt = attempt_no); let fut = op_ref(); let attempt_result = tokio::time::timeout(dur, fut.instrument(span)).await; match attempt_result { @@ -185,7 +198,11 @@ impl ServiceBuilder { return Err(RunError::Inner(e)); } // retryable — backoff and retry - let delay = crate::retry::backoff_delay(&cfg, attempt_no, rand::thread_rng()); + let delay = crate::retry::backoff_delay( + &cfg, + attempt_no, + rand::thread_rng(), + ); tokio::time::sleep(delay).await; } Err(_) => { @@ -194,7 +211,11 @@ impl ServiceBuilder { return Err(RunError::Timeout); } // retryable timeout — backoff and retry - let delay = crate::retry::backoff_delay(&cfg, attempt_no, rand::thread_rng()); + let delay = crate::retry::backoff_delay( + &cfg, + attempt_no, + rand::thread_rng(), + ); tokio::time::sleep(delay).await; } } @@ -202,21 +223,24 @@ impl ServiceBuilder { } else { // retry() expects FnMut() -> Fut (not boxed), so adapt. let mut inner_op = op; - retry(retry_cfg.clone(), |_: &E| true, || { - let span = tracing::info_span!("mytheclipse_service_call"); - let fut = inner_op(); - async move { - fut.instrument(span).await - } - }).await - .map_err(|e| { - self.record(false); - RunError::Retry(e) - }) - .map(|v| { - self.record(true); - v - }) + retry( + retry_cfg.clone(), + |_: &E| true, + || { + let span = tracing::info_span!("mytheclipse_service_call"); + let fut = inner_op(); + async move { fut.instrument(span).await } + }, + ) + .await + .map_err(|e| { + self.record(false); + RunError::Retry(e) + }) + .map(|v| { + self.record(true); + v + }) }; result } else { @@ -224,7 +248,8 @@ impl ServiceBuilder { let mut op = f; let span = tracing::info_span!("mytheclipse_service_call"); let result = if dur > Duration::ZERO { - tokio::time::timeout(dur, op().instrument(span)).await + tokio::time::timeout(dur, op().instrument(span)) + .await .map_err(|_| RunError::Timeout)? .map_err(RunError::Inner) } else { @@ -267,13 +292,19 @@ mod tests { cfg.max_attempts = 3; let builder = ServiceBuilder::new(cfg); let attempts = Arc::new(std::sync::atomic::AtomicU32::new(0)); - let result = builder.run(|| { - let a = Arc::clone(&attempts); - Box::pin(async move { - let n = a.fetch_add(1, std::sync::atomic::Ordering::SeqCst); - if n < 2 { Err::(()) } else { Ok::(42) } + let result = builder + .run(|| { + let a = Arc::clone(&attempts); + Box::pin(async move { + let n = a.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + if n < 2 { + Err::(()) + } else { + Ok::(42) + } + }) }) - }).await; + .await; assert_eq!(result.unwrap(), 42); } @@ -282,10 +313,14 @@ mod tests { let mut cfg = ServiceConfig::default(); cfg.timeout = Duration::from_millis(5); let builder = ServiceBuilder::new(cfg); - let result = builder.run(|| Box::pin(async { - tokio::time::sleep(Duration::from_secs(1)).await; - Ok::<_, ()>(42u32) - })).await; + let result = builder + .run(|| { + Box::pin(async { + tokio::time::sleep(Duration::from_secs(1)).await; + Ok::<_, ()>(42u32) + }) + }) + .await; assert!(matches!(result, Err(RunError::Timeout))); } } diff --git a/crates/mytheclipse/src/shutdown_guard.rs b/crates/mytheclipse/src/shutdown_guard.rs index c257a35..42d647a 100644 --- a/crates/mytheclipse/src/shutdown_guard.rs +++ b/crates/mytheclipse/src/shutdown_guard.rs @@ -37,8 +37,12 @@ use std::sync::{Arc, Mutex}; /// ShutdownGuard::new(move || { d2.fetch_add(1, Ordering::SeqCst); }).finish(); /// assert_eq!(done.load(Ordering::SeqCst), 2); /// ``` +// Type alias for the stored once-only callback, so the nested Arc> field type +// stays within clippy's `type_complexity` threshold. +type GuardFn = Box; + pub struct ShutdownGuard { - inner: Arc>>>, + inner: Arc>>, } impl ShutdownGuard { diff --git a/crates/mytheclipse/tests/race_stress.rs b/crates/mytheclipse/tests/race_stress.rs index eef25bf..fd31080 100644 --- a/crates/mytheclipse/tests/race_stress.rs +++ b/crates/mytheclipse/tests/race_stress.rs @@ -113,7 +113,9 @@ async fn aggregate_error_collects_all_errors_under_contention() { }) .collect(); - let err = AggregateError::from_results(results).err().expect("expected error"); + let err = AggregateError::from_results(results) + .err() + .expect("expected error"); assert_eq!(err.len(), 34, "expected 34 errors (0..99 every 3rd)"); } @@ -157,4 +159,4 @@ async fn concurrency_limiter_max_inflight_never_exceeded() { } assert!(peak.load(Ordering::SeqCst) <= 8, "limiter let >8 inflight"); assert_eq!(inflight.load(Ordering::SeqCst), 0, "limiter leaked permits"); -} \ No newline at end of file +}