docs: doctest examples for round-7-16 features + re-export async_trait
CI / Rustfmt (push) Canceled after 0s
CI / Clippy (push) Canceled after 0s
CI / Test (workspace all features) (push) Canceled after 0s
CI / Test (workspace default features) (push) Canceled after 0s
CI / Test (mytheclipse / bg only) (push) Canceled after 0s
CI / Test (mytheclipse / compute only) (push) Canceled after 0s
CI / Test (mytheclipse / io only) (push) Canceled after 0s
CI / Test (mytheclipse / lifecycle only) (push) Canceled after 0s
CI / Test (mytheclipse / observability only) (push) Canceled after 0s
CI / Test (mytheclipse / resiliency only) (push) Canceled after 0s
CI / Test (mytheclipse / traffic only) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l2-redis) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l1-moka) (push) Canceled after 0s
CI / Test (mytheclipse-cache / default) (push) Canceled after 0s
CI / Test (mytheclipse-config / default) (push) Canceled after 0s
CI / Test (mytheclipse-crypto / default) (push) Canceled after 0s
CI / Test (mytheclipse-event / amqp) (push) Canceled after 0s
CI / Test (mytheclipse-event / nats) (push) Canceled after 0s
CI / Test (mytheclipse-event / default (mem)) (push) Canceled after 0s
CI / Test (mytheclipse-storage / gcs) (push) Canceled after 0s
CI / Test (mytheclipse-storage / s3) (push) Canceled after 0s
CI / Test (mytheclipse-storage / default (local)) (push) Canceled after 0s
CI / Run mytheclipse example (push) Canceled after 0s
CI / Docs check (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-cache) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-config) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-crypto) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-event) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-storage) (push) Canceled after 0s
Release / Semantic Release (push) Canceled after 0s
CI / Rustfmt (push) Canceled after 0s
CI / Clippy (push) Canceled after 0s
CI / Test (workspace all features) (push) Canceled after 0s
CI / Test (workspace default features) (push) Canceled after 0s
CI / Test (mytheclipse / bg only) (push) Canceled after 0s
CI / Test (mytheclipse / compute only) (push) Canceled after 0s
CI / Test (mytheclipse / io only) (push) Canceled after 0s
CI / Test (mytheclipse / lifecycle only) (push) Canceled after 0s
CI / Test (mytheclipse / observability only) (push) Canceled after 0s
CI / Test (mytheclipse / resiliency only) (push) Canceled after 0s
CI / Test (mytheclipse / traffic only) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l2-redis) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l1-moka) (push) Canceled after 0s
CI / Test (mytheclipse-cache / default) (push) Canceled after 0s
CI / Test (mytheclipse-config / default) (push) Canceled after 0s
CI / Test (mytheclipse-crypto / default) (push) Canceled after 0s
CI / Test (mytheclipse-event / amqp) (push) Canceled after 0s
CI / Test (mytheclipse-event / nats) (push) Canceled after 0s
CI / Test (mytheclipse-event / default (mem)) (push) Canceled after 0s
CI / Test (mytheclipse-storage / gcs) (push) Canceled after 0s
CI / Test (mytheclipse-storage / s3) (push) Canceled after 0s
CI / Test (mytheclipse-storage / default (local)) (push) Canceled after 0s
CI / Run mytheclipse example (push) Canceled after 0s
CI / Docs check (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-cache) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-config) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-crypto) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-event) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-storage) (push) Canceled after 0s
Release / Semantic Release (push) Canceled after 0s
This commit is contained in:
@@ -10,6 +10,27 @@
|
|||||||
use std::fmt;
|
use std::fmt;
|
||||||
|
|
||||||
/// An error that groups one or more underlying errors.
|
/// 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)]
|
#[derive(Debug)]
|
||||||
pub struct AggregateError {
|
pub struct AggregateError {
|
||||||
errors: Vec<Box<dyn std::error::Error + Send + Sync>>,
|
errors: Vec<Box<dyn std::error::Error + Send + Sync>>,
|
||||||
|
|||||||
@@ -19,6 +19,22 @@ use crate::metrics::MetricsCollector;
|
|||||||
use crate::service_builder::{RunError, ServiceBuilder, ServiceConfig};
|
use crate::service_builder::{RunError, ServiceBuilder, ServiceConfig};
|
||||||
|
|
||||||
/// A [`ServiceBuilder`] wrapper that auto-records latency and outcome metrics.
|
/// 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 {
|
pub struct AutoMetricsServiceBuilder {
|
||||||
inner: ServiceBuilder,
|
inner: ServiceBuilder,
|
||||||
metrics: MetricsCollector,
|
metrics: MetricsCollector,
|
||||||
|
|||||||
@@ -128,6 +128,10 @@ pub use backpressure::{BackpressureError, BackpressureQueue, OverflowPolicy};
|
|||||||
pub use concurrency::{ConcurrencyLimiter, ConcurrencyPermit};
|
pub use concurrency::{ConcurrencyLimiter, ConcurrencyPermit};
|
||||||
#[cfg(feature = "traffic")]
|
#[cfg(feature = "traffic")]
|
||||||
pub use pool::{Pool, PoolError, Pooled, SemaphorePool, AutoReconnectPool, Reconnectable};
|
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")]
|
#[cfg(feature = "lifecycle")]
|
||||||
pub use shutdown::{ShutdownManager, ShutdownSignal};
|
pub use shutdown::{ShutdownManager, ShutdownSignal};
|
||||||
|
|||||||
@@ -22,6 +22,19 @@ use crate::aggregate_error::AggregateError;
|
|||||||
///
|
///
|
||||||
/// Note: `items` is fully collected into memory up front (see
|
/// Note: `items` is fully collected into memory up front (see
|
||||||
/// [`parallel_for_each`] for a streaming variant that avoids materializing).
|
/// [`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<I, T, F, Fut, E>(
|
pub async fn parallel_map<I, T, F, Fut, E>(
|
||||||
items: I,
|
items: I,
|
||||||
concurrency: usize,
|
concurrency: usize,
|
||||||
@@ -125,6 +138,28 @@ where
|
|||||||
/// task and never gets ahead.
|
/// task and never gets ahead.
|
||||||
///
|
///
|
||||||
/// Errors are aggregated into a single [`AggregateError`].
|
/// 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<I, F, Fut, E>(
|
pub async fn parallel_for_each<I, F, Fut, E>(
|
||||||
items: I,
|
items: I,
|
||||||
concurrency: usize,
|
concurrency: usize,
|
||||||
|
|||||||
@@ -94,6 +94,36 @@ pub trait Reconnectable {
|
|||||||
/// [`Reconnectable::is_healthy`]; if unhealthy, a replacement is produced via
|
/// [`Reconnectable::is_healthy`]; if unhealthy, a replacement is produced via
|
||||||
/// [`Reconnectable::reconnect`] and handed back instead. This removes the
|
/// [`Reconnectable::reconnect`] and handed back instead. This removes the
|
||||||
/// per-call-site "is my connection dead? rebuild it" boilerplate.
|
/// 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<Conn, Box<dyn std::error::Error + Send + Sync>> {
|
||||||
|
/// 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<P, R> {
|
pub struct AutoReconnectPool<P, R> {
|
||||||
inner: P,
|
inner: P,
|
||||||
reconnect: R,
|
reconnect: R,
|
||||||
|
|||||||
@@ -9,6 +9,36 @@ use std::pin::Pin;
|
|||||||
use crate::retry::{retry, RetryConfig, RetryError};
|
use crate::retry::{retry, RetryConfig, RetryError};
|
||||||
|
|
||||||
/// Extension trait adding ergonomic `.retry()` to any fallible future.
|
/// 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<T, E>: Future<Output = Result<T, E>> + Sized + 'static
|
pub trait RetryExt<T, E>: Future<Output = Result<T, E>> + Sized + 'static
|
||||||
where
|
where
|
||||||
E: std::fmt::Debug + Send + 'static,
|
E: std::fmt::Debug + Send + 'static,
|
||||||
|
|||||||
@@ -20,6 +20,22 @@ use std::num::NonZeroUsize;
|
|||||||
///
|
///
|
||||||
/// Each field is a concrete `usize` (never zero) so it can be fed directly
|
/// 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.
|
/// 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)]
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
pub struct RuntimeConfig {
|
pub struct RuntimeConfig {
|
||||||
/// Worker threads for the main async runtime (default: one per core).
|
/// Worker threads for the main async runtime (default: one per core).
|
||||||
|
|||||||
@@ -18,6 +18,25 @@ use std::sync::{Arc, Mutex};
|
|||||||
/// The callback is wrapped in a `Mutex<Option<_>>` so it can be taken out and
|
/// The callback is wrapped in a `Mutex<Option<_>>` so it can be taken out and
|
||||||
/// run exactly once — even on a `panic!`-unwound drop — guaranteeing at-most-
|
/// run exactly once — even on a `panic!`-unwound drop — guaranteeing at-most-
|
||||||
/// once semantics (no double-shutdown race).
|
/// 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 {
|
pub struct ShutdownGuard {
|
||||||
inner: Arc<Mutex<Option<Box<dyn FnOnce() + Send>>>>,
|
inner: Arc<Mutex<Option<Box<dyn FnOnce() + Send>>>>,
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user