From a6cf58ea7fe775f1e7ed4781d31435459f440974 Mon Sep 17 00:00:00 2001 From: asepharyana Date: Sat, 29 Aug 2026 21:47:20 +0700 Subject: [PATCH] docs: doctest examples for round-7-16 features + re-export async_trait --- crates/mytheclipse/src/aggregate_error.rs | 21 +++++++++++ .../mytheclipse/src/auto_metrics_service.rs | 16 +++++++++ crates/mytheclipse/src/lib.rs | 4 +++ crates/mytheclipse/src/parallel_map.rs | 35 +++++++++++++++++++ crates/mytheclipse/src/pool.rs | 30 ++++++++++++++++ crates/mytheclipse/src/retry_ext.rs | 30 ++++++++++++++++ crates/mytheclipse/src/runtime_auto.rs | 16 +++++++++ crates/mytheclipse/src/shutdown_guard.rs | 19 ++++++++++ 8 files changed, 171 insertions(+) diff --git a/crates/mytheclipse/src/aggregate_error.rs b/crates/mytheclipse/src/aggregate_error.rs index e1e4ede..86f0061 100644 --- a/crates/mytheclipse/src/aggregate_error.rs +++ b/crates/mytheclipse/src/aggregate_error.rs @@ -10,6 +10,27 @@ use std::fmt; /// An error that groups one or more underlying errors. +/// +/// ``` +/// use mytheclipse::aggregate_error::AggregateError; +/// +/// // Collect errors from N parallel results — all of them, not just the first: +/// let results = vec![ +/// Ok::<_, std::io::Error>(1), +/// Err(std::io::Error::other("boom")), +/// Ok::<_, std::io::Error>(3), +/// Err(std::io::Error::other("bam")), +/// ]; +/// let out = AggregateError::from_results(results); +/// let err = out.unwrap_err(); +/// assert_eq!(err.len(), 2); // both errors collected +/// +/// // Build one incrementally too: +/// let mut agg = AggregateError::with_context("fan-out"); +/// agg.push(std::io::Error::other("first")); +/// agg.push(std::io::Error::other("second")); +/// assert_eq!(agg.len(), 2); +/// ``` #[derive(Debug)] pub struct AggregateError { errors: Vec>, diff --git a/crates/mytheclipse/src/auto_metrics_service.rs b/crates/mytheclipse/src/auto_metrics_service.rs index 7561554..c5f0ab8 100644 --- a/crates/mytheclipse/src/auto_metrics_service.rs +++ b/crates/mytheclipse/src/auto_metrics_service.rs @@ -19,6 +19,22 @@ use crate::metrics::MetricsCollector; use crate::service_builder::{RunError, ServiceBuilder, ServiceConfig}; /// A [`ServiceBuilder`] wrapper that auto-records latency and outcome metrics. +/// +/// ``` +/// use std::time::Duration; +/// use mytheclipse::auto_metrics_service::AutoMetricsServiceBuilder; +/// use mytheclipse::service_builder::ServiceConfig; +/// +/// let cfg = ServiceConfig { +/// max_attempts: 2, +/// timeout: Duration::from_millis(1), +/// ..ServiceConfig::default() +/// }; +/// let svc = AutoMetricsServiceBuilder::new("checkout", cfg); +/// // Every `.run()` call now auto-records outcome + latency on the shared +/// // MetricsCollector (retrievable via `collector()`). +/// let _ = svc.collector(); +/// ``` pub struct AutoMetricsServiceBuilder { inner: ServiceBuilder, metrics: MetricsCollector, diff --git a/crates/mytheclipse/src/lib.rs b/crates/mytheclipse/src/lib.rs index 5370f70..1c14003 100644 --- a/crates/mytheclipse/src/lib.rs +++ b/crates/mytheclipse/src/lib.rs @@ -128,6 +128,10 @@ pub use backpressure::{BackpressureError, BackpressureQueue, OverflowPolicy}; 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; #[cfg(feature = "lifecycle")] pub use shutdown::{ShutdownManager, ShutdownSignal}; diff --git a/crates/mytheclipse/src/parallel_map.rs b/crates/mytheclipse/src/parallel_map.rs index 1e4fdca..e06431b 100644 --- a/crates/mytheclipse/src/parallel_map.rs +++ b/crates/mytheclipse/src/parallel_map.rs @@ -22,6 +22,19 @@ use crate::aggregate_error::AggregateError; /// /// Note: `items` is fully collected into memory up front (see /// [`parallel_for_each`] for a streaming variant that avoids materializing). +/// +/// ``` +/// use mytheclipse::parallel_map::parallel_map; +/// +/// #[tokio::main] +/// async fn main() { +/// let items = vec![1u32, 2, 3, 4, 5]; +/// let doubled = parallel_map(items, 4, |x| async move { Ok::<_, std::io::Error>(x * 2) }) +/// .await +/// .unwrap(); +/// assert_eq!(doubled, vec![2, 4, 6, 8, 10]); +/// } +/// ``` pub async fn parallel_map( items: I, concurrency: usize, @@ -125,6 +138,28 @@ where /// task and never gets ahead. /// /// Errors are aggregated into a single [`AggregateError`]. +/// +/// ``` +/// use std::sync::atomic::{AtomicUsize, Ordering}; +/// use std::sync::Arc; +/// use mytheclipse::parallel_map::parallel_for_each; +/// +/// #[tokio::main] +/// async fn main() { +/// let seen = Arc::new(AtomicUsize::new(0)); +/// let s = Arc::clone(&seen); +/// parallel_for_each(0u32..100, 8, move |x| { +/// let s = Arc::clone(&s); +/// async move { +/// s.fetch_add(x as usize, Ordering::SeqCst); +/// Ok::<_, std::io::Error>(()) +/// } +/// }) +/// .await +/// .unwrap(); +/// assert_eq!(seen.load(Ordering::SeqCst), 4950); // sum 0..100 +/// } +/// ``` pub async fn parallel_for_each( items: I, concurrency: usize, diff --git a/crates/mytheclipse/src/pool.rs b/crates/mytheclipse/src/pool.rs index bd1bfa8..804f2b0 100644 --- a/crates/mytheclipse/src/pool.rs +++ b/crates/mytheclipse/src/pool.rs @@ -94,6 +94,36 @@ pub trait Reconnectable { /// [`Reconnectable::is_healthy`]; if unhealthy, a replacement is produced via /// [`Reconnectable::reconnect`] and handed back instead. This removes the /// per-call-site "is my connection dead? rebuild it" boilerplate. +/// +/// ``` +/// use mytheclipse::pool::{AutoReconnectPool, Pool, Reconnectable, SemaphorePool}; +/// use mytheclipse::async_trait; +/// +/// // A "connection" that's dead when its value equals 0. +/// #[derive(Clone)] +/// struct Conn { alive: bool } +/// impl Default for Conn { fn default() -> Self { Self { alive: true } } } +/// +/// #[derive(Default)] +/// struct ConnReconnector; +/// +/// #[async_trait] +/// impl Reconnectable for ConnReconnector { +/// type Item = Conn; +/// fn is_healthy(&self, c: &Conn) -> bool { c.alive } +/// async fn reconnect(&self) -> Result> { +/// Ok(Conn::default()) +/// } +/// } +/// +/// #[tokio::main] +/// async fn main() { +/// let pool = SemaphorePool::new(vec![Conn { alive: false }, Conn::default()]); +/// let pool = AutoReconnectPool::new(pool, ConnReconnector); +/// let first = pool.acquire().await.unwrap(); +/// assert!(first.resource.alive); // dead one was transparently replaced +/// } +/// ``` pub struct AutoReconnectPool { inner: P, reconnect: R, diff --git a/crates/mytheclipse/src/retry_ext.rs b/crates/mytheclipse/src/retry_ext.rs index f5cc7c8..8c173da 100644 --- a/crates/mytheclipse/src/retry_ext.rs +++ b/crates/mytheclipse/src/retry_ext.rs @@ -9,6 +9,36 @@ use std::pin::Pin; use crate::retry::{retry, RetryConfig, RetryError}; /// Extension trait adding ergonomic `.retry()` to any fallible future. +/// +/// ``` +/// use std::time::Duration; +/// use mytheclipse::retry_ext::RetryExt; +/// use mytheclipse::retry::RetryConfig; +/// +/// #[tokio::main] +/// async fn main() { +/// let cfg = RetryConfig { max_attempts: 3, base_delay: Duration::from_millis(1), ..RetryConfig::default() }; +/// let attempts = std::sync::Arc::new(std::sync::atomic::AtomicU32::new(0)); +/// let a = std::sync::Arc::clone(&attempts); +/// +/// let op = move || { +/// let a = std::sync::Arc::clone(&a); +/// async move { +/// if a.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2 { +/// Err::<(), String>("transient".into()) +/// } else { +/// Ok(()) +/// } +/// } +/// }; +/// +/// let result = async { Err::<(), String>("first".into()) } +/// .retry(cfg, |_: &String| true, op) +/// .await; +/// assert!(result.is_ok()); +/// assert_eq!(attempts.load(std::sync::atomic::Ordering::SeqCst), 3); +/// } +/// ``` pub trait RetryExt: Future> + Sized + 'static where E: std::fmt::Debug + Send + 'static, diff --git a/crates/mytheclipse/src/runtime_auto.rs b/crates/mytheclipse/src/runtime_auto.rs index ddea339..4230c91 100644 --- a/crates/mytheclipse/src/runtime_auto.rs +++ b/crates/mytheclipse/src/runtime_auto.rs @@ -20,6 +20,22 @@ use std::num::NonZeroUsize; /// /// Each field is a concrete `usize` (never zero) so it can be fed directly /// into a tokio/rayon/std-thread builder with no further `max(1)` guards. +/// +/// ``` +/// use mytheclipse::runtime_auto::RuntimeConfig; +/// +/// let cfg = RuntimeConfig::auto(); // from the host CPU +/// let sized = RuntimeConfig::from_cores(4); // or explicit +/// +/// // Feed straight into a tokio builder: +/// let rt = tokio::runtime::Builder::new_multi_thread() +/// .worker_threads(cfg.worker_threads) +/// .max_blocking_threads(cfg.max_blocking_threads) +/// .build() +/// .unwrap(); +/// let _ = sized.compute_threads; +/// rt.shutdown_timeout(std::time::Duration::from_millis(1)); +/// ``` #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub struct RuntimeConfig { /// Worker threads for the main async runtime (default: one per core). diff --git a/crates/mytheclipse/src/shutdown_guard.rs b/crates/mytheclipse/src/shutdown_guard.rs index 35cd175..c257a35 100644 --- a/crates/mytheclipse/src/shutdown_guard.rs +++ b/crates/mytheclipse/src/shutdown_guard.rs @@ -18,6 +18,25 @@ use std::sync::{Arc, Mutex}; /// The callback is wrapped in a `Mutex>` so it can be taken out and /// run exactly once — even on a `panic!`-unwound drop — guaranteeing at-most- /// once semantics (no double-shutdown race). +/// +/// ``` +/// use std::sync::atomic::{AtomicUsize, Ordering}; +/// use std::sync::Arc; +/// use mytheclipse::shutdown_guard::ShutdownGuard; +/// +/// let done = Arc::new(AtomicUsize::new(0)); +/// let d = Arc::clone(&done); +/// { +/// let _guard = ShutdownGuard::new(move || { d.fetch_add(1, Ordering::SeqCst); }); +/// // do work... +/// } // guard dropped here -> callback fires exactly once +/// assert_eq!(done.load(Ordering::SeqCst), 1); +/// +/// // Or fire early with `finish()` (disarms the drop): +/// let d2 = Arc::clone(&done); +/// ShutdownGuard::new(move || { d2.fetch_add(1, Ordering::SeqCst); }).finish(); +/// assert_eq!(done.load(Ordering::SeqCst), 2); +/// ``` pub struct ShutdownGuard { inner: Arc>>>, }