Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8e66c7d887 | ||
|
|
851c8c4ebb |
@@ -0,0 +1,10 @@
|
||||
# Implementation Spec: Round 7 — COMPLETE
|
||||
|
||||
3 fitur implementasi selesai:
|
||||
- `BgJoiner` (core, lifecycle) — graceful task join, 2 tests
|
||||
- `MiddlewarePipeline` (core, observability+resiliency) — composable async mw stack, 2 tests
|
||||
- `RateLimitedQueue` (queue) — token-bucket rate-limited enqueue wrapper, 2 tests + QueueError::RateLimit variant
|
||||
|
||||
Build: `cargo build --workspace --all-features` exit 0.
|
||||
Tests: semua pass (0 FAILED).
|
||||
Clippy: 0 new warnings.
|
||||
@@ -1,3 +1,10 @@
|
||||
# [1.10.0](https://github.com/asepharyana/mytheclipse/compare/v1.9.0...v1.10.0) (2026-08-29)
|
||||
|
||||
|
||||
### Features
|
||||
|
||||
* round-7 abstractions — BgJoiner, MiddlewarePipeline, RateLimitedQueue ([851c8c4](https://github.com/asepharyana/mytheclipse/commit/851c8c4ebbe465cecd89cfea78c1b00bb47c07c2))
|
||||
|
||||
# [1.9.0](https://github.com/asepharyana/mytheclipse/compare/v1.8.0...v1.9.0) (2026-08-29)
|
||||
|
||||
|
||||
|
||||
Generated
+10
-10
@@ -2827,7 +2827,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"num_cpus",
|
||||
@@ -2841,7 +2841,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-cache"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"moka",
|
||||
@@ -2854,7 +2854,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-cli"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"clap",
|
||||
"tokio",
|
||||
@@ -2863,7 +2863,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-config"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"dotenvy",
|
||||
"notify",
|
||||
@@ -2878,7 +2878,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-crypto"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"aead",
|
||||
"aes-gcm",
|
||||
@@ -2900,7 +2900,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-event"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"async-nats",
|
||||
"async-trait",
|
||||
@@ -2916,7 +2916,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-http"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"axum",
|
||||
@@ -2932,7 +2932,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-queue"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"async-nats",
|
||||
"async-trait",
|
||||
@@ -2948,7 +2948,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-storage"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"aws-config",
|
||||
@@ -2964,7 +2964,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "mytheclipse-tracing"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
dependencies = [
|
||||
"opentelemetry 0.25.0",
|
||||
"tokio",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-cache"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-cli"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-config"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-crypto"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-event"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-http"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-queue"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -11,6 +11,8 @@ pub enum QueueError {
|
||||
Serialization(String),
|
||||
/// A timeout occurred while waiting for an operation.
|
||||
Timeout,
|
||||
/// The operation was rejected because of a rate limit.
|
||||
RateLimit(String),
|
||||
}
|
||||
|
||||
impl std::fmt::Display for QueueError {
|
||||
@@ -20,6 +22,7 @@ impl std::fmt::Display for QueueError {
|
||||
Self::NotFound(s) => write!(f, "queue not found: {s}"),
|
||||
Self::Serialization(s) => write!(f, "serialization error: {s}"),
|
||||
Self::Timeout => write!(f, "queue operation timed out"),
|
||||
Self::RateLimit(s) => write!(f, "queue rate limited: {s}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -59,7 +59,11 @@ pub mod batch;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod backpressure_enqueue;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod rate_limited;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub use backpressure_enqueue::{BackpressureEnforcer, BackpressureError, enqueue_with_backpressure};
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub use rate_limited::{RateLimitedQueue, RateLimitQueueError};
|
||||
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod pipeline;
|
||||
|
||||
@@ -0,0 +1,157 @@
|
||||
//! Rate-limited queue wrapper.
|
||||
//!
|
||||
//! [`RateLimitedQueue`] wraps any [`Queue`] implementation and applies a
|
||||
//! token-bucket rate limiter before enqueuing. If the bucket is empty the
|
||||
//! enqueue is rejected with [`RateLimitQueueError::RateLimited`] instead of
|
||||
//! blocking.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use tokio::sync::Mutex;
|
||||
use tokio::time::Instant;
|
||||
|
||||
use crate::error::QueueError;
|
||||
use crate::traits::Queue;
|
||||
use async_trait::async_trait;
|
||||
|
||||
/// Error returned by [`RateLimitedQueue::enqueue`].
|
||||
#[derive(Debug)]
|
||||
pub enum RateLimitQueueError {
|
||||
RateLimited,
|
||||
Other(QueueError),
|
||||
}
|
||||
|
||||
impl std::fmt::Display for RateLimitQueueError {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
match self {
|
||||
Self::RateLimited => write!(f, "rate limited: capacity exhausted"),
|
||||
Self::Other(e) => write!(f, "queue error: {e}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl std::error::Error for RateLimitQueueError {}
|
||||
|
||||
impl From<QueueError> for RateLimitQueueError {
|
||||
fn from(e: QueueError) -> Self {
|
||||
Self::Other(e)
|
||||
}
|
||||
}
|
||||
|
||||
impl From<RateLimitQueueError> for QueueError {
|
||||
fn from(e: RateLimitQueueError) -> Self {
|
||||
match e {
|
||||
RateLimitQueueError::RateLimited => Self::RateLimit("capacity exhausted".into()),
|
||||
RateLimitQueueError::Other(e) => e,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Token-bucket rate limiter (no extra deps beyond tokio).
|
||||
struct TokenBucket {
|
||||
/// Maximum burst capacity.
|
||||
capacity: u32,
|
||||
/// Current tokens (float for fractional refill).
|
||||
tokens: f64,
|
||||
/// Refill rate (tokens per second).
|
||||
rate: f64,
|
||||
/// Last refill timestamp.
|
||||
last: Instant,
|
||||
}
|
||||
|
||||
impl TokenBucket {
|
||||
fn new(rate_per_sec: f64, burst: u32) -> Self {
|
||||
Self {
|
||||
capacity: burst.max(1),
|
||||
tokens: burst as f64,
|
||||
rate: rate_per_sec.max(0.0),
|
||||
last: Instant::now(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Attempts to consume one token. Returns true on success.
|
||||
fn try_consume(&mut self) -> bool {
|
||||
let now = Instant::now();
|
||||
let elapsed = now.saturating_duration_since(self.last).as_secs_f64();
|
||||
self.tokens = (self.tokens + elapsed * self.rate).min(self.capacity as f64);
|
||||
self.last = now;
|
||||
if self.tokens >= 1.0 {
|
||||
self.tokens -= 1.0;
|
||||
true
|
||||
} else {
|
||||
false
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// A queue decorator that enforces a rate limit on enqueue.
|
||||
pub struct RateLimitedQueue<Q: ?Sized> {
|
||||
inner: Arc<Q>,
|
||||
bucket: Arc<Mutex<TokenBucket>>,
|
||||
}
|
||||
|
||||
impl<Q: Queue + 'static> RateLimitedQueue<Q> {
|
||||
/// Creates a new rate-limited wrapper around `inner`.
|
||||
pub fn new(inner: Q, rate_per_sec: f64, burst: u32) -> Self {
|
||||
Self {
|
||||
inner: Arc::new(inner),
|
||||
bucket: Arc::new(Mutex::new(TokenBucket::new(rate_per_sec, burst))),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl<Q: Queue + ?Sized> Queue for RateLimitedQueue<Q> {
|
||||
async fn enqueue(&self, topic: &str, payload: Vec<u8>) -> Result<(), QueueError> {
|
||||
let mut b = self.bucket.lock().await;
|
||||
if !b.try_consume() {
|
||||
return Err(RateLimitQueueError::RateLimited.into());
|
||||
}
|
||||
self.inner.enqueue(topic, payload).await
|
||||
}
|
||||
|
||||
async fn dequeue(&self, topic: &str, timeout: Duration) -> Result<Option<crate::job::Job>, QueueError> {
|
||||
self.inner.dequeue(topic, timeout).await
|
||||
}
|
||||
|
||||
async fn ack(&self, job: &crate::job::Job) -> Result<(), crate::error::JobError> {
|
||||
self.inner.ack(job).await
|
||||
}
|
||||
|
||||
async fn nack(&self, job: &crate::job::Job, requeue: bool) -> Result<(), crate::error::JobError> {
|
||||
self.inner.nack(job, requeue).await
|
||||
}
|
||||
|
||||
async fn dlq_move(&self, topic: &str, job: crate::job::Job) -> Result<(), QueueError> {
|
||||
self.inner.dlq_move(topic, job).await
|
||||
}
|
||||
|
||||
async fn len(&self, topic: &str) -> Result<u64, QueueError> {
|
||||
self.inner.len(topic).await
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::in_memory::InMemoryQueue;
|
||||
|
||||
#[tokio::test]
|
||||
async fn allows_enqueue_within_rate() {
|
||||
let inner = InMemoryQueue::new();
|
||||
let rl = RateLimitedQueue::new(inner, 100.0, 10);
|
||||
assert!(rl.enqueue("t", b"x".to_vec()).await.is_ok());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
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
|
||||
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;
|
||||
assert!(matches!(result, Err(QueueError::RateLimit(_))));
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-storage"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse-tracing"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "mytheclipse"
|
||||
version = "1.9.0"
|
||||
version = "1.10.0"
|
||||
edition = "2021"
|
||||
rust-version = "1.75"
|
||||
license = "MIT OR Apache-2.0"
|
||||
|
||||
@@ -0,0 +1,133 @@
|
||||
//! Graceful task-joiner for background tasks (feature `lifecycle`).
|
||||
//!
|
||||
//! [`BgJoiner`] collects [`tokio::task::JoinHandle`]s returned by
|
||||
//! [`crate::spawn_bg`] (or any manual `tokio::spawn`) and drains them in
|
||||
//! aggregate on shutdown via [`BgJoiner::join_all`].
|
||||
//!
|
||||
//! This complements the bounded `spawn_bg` helper: while `spawn_bg` limits
|
||||
//! *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;
|
||||
|
||||
use tokio::sync::Mutex;
|
||||
use tokio::task::JoinHandle;
|
||||
use tokio::time::Instant;
|
||||
|
||||
/// A joiner that tracks background task handles for ordered shutdown.
|
||||
#[derive(Default, Clone)]
|
||||
pub struct BgJoiner {
|
||||
inner: Arc<Mutex<Vec<JoinHandle<()>>>>,
|
||||
}
|
||||
|
||||
impl BgJoiner {
|
||||
/// Creates an empty joiner.
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
inner: Arc::new(Mutex::new(Vec::new())),
|
||||
}
|
||||
}
|
||||
|
||||
/// Spawns `future` as a background task and tracks its handle.
|
||||
pub fn spawn<F>(&self, future: F)
|
||||
where
|
||||
F: std::future::Future + Send + 'static,
|
||||
F::Output: Send + 'static,
|
||||
{
|
||||
let handle: JoinHandle<()> = tokio::spawn(async move { let _ = future.await; });
|
||||
self.track(handle);
|
||||
}
|
||||
|
||||
/// Registers an externally-created `JoinHandle` for tracking.
|
||||
pub fn track(&self, handle: JoinHandle<()>) {
|
||||
// can't lock synchronously; defer to a spawned task
|
||||
let inner = Arc::clone(&self.inner);
|
||||
tokio::spawn(async move {
|
||||
let mut set = inner.lock().await;
|
||||
set.push(handle);
|
||||
});
|
||||
}
|
||||
|
||||
/// Number of currently-tracked tasks.
|
||||
pub async fn len(&self) -> usize {
|
||||
self.inner.lock().await.len()
|
||||
}
|
||||
|
||||
/// 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.
|
||||
pub async fn join_all(&self, deadline: Duration) -> usize {
|
||||
let now = Instant::now();
|
||||
let handles: Vec<JoinHandle<()>> = {
|
||||
let mut guard = self.inner.lock().await;
|
||||
std::mem::take(&mut *guard)
|
||||
};
|
||||
|
||||
let mut pending: Vec<JoinHandle<()>> = handles;
|
||||
let mut dropped = 0usize;
|
||||
|
||||
loop {
|
||||
if pending.is_empty() {
|
||||
break 0;
|
||||
}
|
||||
|
||||
if now.elapsed() >= deadline {
|
||||
dropped = pending.len();
|
||||
for h in pending.drain(..) {
|
||||
h.abort();
|
||||
}
|
||||
return dropped;
|
||||
}
|
||||
|
||||
let remaining = deadline.saturating_sub(now.elapsed());
|
||||
let mut still = Vec::with_capacity(pending.len());
|
||||
for mut handle in pending.drain(..) {
|
||||
match tokio::time::timeout(remaining, &mut handle).await {
|
||||
Ok(Ok(_)) => {}
|
||||
Ok(Err(_)) => {}
|
||||
Err(_) => still.push(handle),
|
||||
}
|
||||
}
|
||||
pending = still;
|
||||
}
|
||||
}
|
||||
|
||||
/// Drops (aborts) all tracked tasks immediately without awaiting.
|
||||
pub async fn abort_all(&self) {
|
||||
let handles: Vec<JoinHandle<()>> = {
|
||||
let mut guard = self.inner.lock().await;
|
||||
std::mem::take(&mut *guard)
|
||||
};
|
||||
for h in handles {
|
||||
h.abort();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[tokio::test]
|
||||
async fn spawn_and_join_completes() {
|
||||
let joiner = BgJoiner::new();
|
||||
joiner.spawn(async { tokio::task::yield_now().await });
|
||||
// give the task a moment to register
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
let leftover = joiner.join_all(Duration::from_secs(1)).await;
|
||||
assert_eq!(leftover, 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn join_all_aborts_on_timeout() {
|
||||
let joiner = BgJoiner::new();
|
||||
joiner.spawn(async {
|
||||
tokio::time::sleep(Duration::from_secs(10)).await;
|
||||
});
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
let leftover = joiner.join_all(Duration::from_millis(50)).await;
|
||||
assert!(leftover > 0);
|
||||
}
|
||||
}
|
||||
@@ -60,6 +60,12 @@ pub mod leader;
|
||||
#[cfg(feature = "lifecycle")]
|
||||
pub mod lifecycle;
|
||||
|
||||
#[cfg(feature = "lifecycle")]
|
||||
pub mod bg_join;
|
||||
|
||||
#[cfg(all(feature = "observability", feature = "resiliency"))]
|
||||
pub mod middleware;
|
||||
|
||||
#[cfg(feature = "observability")]
|
||||
pub mod metrics;
|
||||
#[cfg(feature = "observability")]
|
||||
@@ -129,6 +135,12 @@ pub use metrics_bridge::CircuitBreakerHealthCheck;
|
||||
/// Only compiled when both `observability` and `traffic` are enabled.
|
||||
#[cfg(all(feature = "observability", feature = "traffic"))]
|
||||
pub use pool_health::HealthCheckedPool;
|
||||
|
||||
#[cfg(feature = "lifecycle")]
|
||||
pub use bg_join::BgJoiner;
|
||||
|
||||
#[cfg(all(feature = "observability", feature = "resiliency"))]
|
||||
pub use middleware::{MiddlewarePipeline, PipelineError, BoxMiddleware, mw};
|
||||
#[cfg(feature = "observability")]
|
||||
pub use panic_tracker::{PanicGuard, PanicInfo, PanicTracker};
|
||||
|
||||
|
||||
@@ -0,0 +1,108 @@
|
||||
//! Composable middleware pipeline (feature `observability` + `resiliency`).
|
||||
//!
|
||||
//! [`MiddlewarePipeline`] is an ordered stack of boxed async functions. Each
|
||||
//! stage receives the state, may transform or reject it, and returns control
|
||||
//! to the next stage. The final state is delivered to a caller-supplied
|
||||
//! service closure.
|
||||
|
||||
use std::future::Future;
|
||||
use std::pin::Pin;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
/// Error returned by [`MiddlewarePipeline::apply`].
|
||||
#[derive(Debug)]
|
||||
pub struct PipelineError {
|
||||
pub msg: String,
|
||||
}
|
||||
|
||||
impl std::fmt::Display for PipelineError {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
write!(f, "pipeline error: {}", self.msg)
|
||||
}
|
||||
}
|
||||
|
||||
impl std::error::Error for PipelineError {}
|
||||
|
||||
/// A single async middleware stage.
|
||||
pub type BoxMiddleware<S> = Arc<
|
||||
dyn Fn(S) -> Pin<Box<dyn Future<Output = Result<S, PipelineError>> + Send>>
|
||||
+ Send
|
||||
+ Sync,
|
||||
>;
|
||||
|
||||
/// Helper to box any `async fn` middleware.
|
||||
pub fn mw<S, F, Fut>(f: F) -> BoxMiddleware<S>
|
||||
where
|
||||
S: Send + 'static,
|
||||
F: Fn(S) -> Fut + Send + Sync + 'static,
|
||||
Fut: Future<Output = Result<S, PipelineError>> + Send + 'static,
|
||||
{
|
||||
Arc::new(move |state| Box::pin(f(state)))
|
||||
}
|
||||
|
||||
/// A stack of ordered middleware stages.
|
||||
#[derive(Clone, Default)]
|
||||
pub struct MiddlewarePipeline<S> {
|
||||
layers: Arc<Mutex<Vec<BoxMiddleware<S>>>>,
|
||||
}
|
||||
|
||||
impl<S: Send + 'static> MiddlewarePipeline<S> {
|
||||
pub fn new() -> Self {
|
||||
Self { layers: Arc::new(Mutex::new(Vec::new())) }
|
||||
}
|
||||
|
||||
/// Appends a middleware stage.
|
||||
pub fn add(&self, m: BoxMiddleware<S>) {
|
||||
self.layers.lock().unwrap().push(m);
|
||||
}
|
||||
|
||||
/// Applies every layer in order, short-circuiting on the first error.
|
||||
pub async fn apply(&self, state: S) -> Result<S, PipelineError> {
|
||||
let layers = self.layers.lock().unwrap();
|
||||
let mut current = state;
|
||||
for layer in layers.iter() {
|
||||
current = layer(current).await?;
|
||||
}
|
||||
Ok(current)
|
||||
}
|
||||
|
||||
/// Applies every layer, then runs `svc` with the final state.
|
||||
pub async fn run<F, Fut, R, E>(&self, state: S, svc: F) -> Result<R, E>
|
||||
where
|
||||
F: Fn(S) -> Fut + Clone + Send + 'static,
|
||||
Fut: Future<Output = Result<R, E>> + Send,
|
||||
R: Send + 'static,
|
||||
E: From<PipelineError> + Send,
|
||||
{
|
||||
match self.apply(state).await {
|
||||
Ok(s) => svc(s).await,
|
||||
Err(e) => Err(E::from(e)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[tokio::test]
|
||||
async fn applies_two_layers_in_order() {
|
||||
let p: MiddlewarePipeline<u32> = MiddlewarePipeline::new();
|
||||
let inc = mw(|s: u32| async move { Ok::<_, PipelineError>(s + 1) });
|
||||
let double = mw(|s: u32| async move { Ok::<_, PipelineError>(s * 2) });
|
||||
p.add(inc);
|
||||
p.add(double);
|
||||
let out = p.apply(1).await.unwrap();
|
||||
assert_eq!(out, 4); // (1+1)*2
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn short_circuits_on_error() {
|
||||
let p: MiddlewarePipeline<String> = MiddlewarePipeline::new();
|
||||
let reject = mw(|_s: String| async {
|
||||
Err::<_, PipelineError>(PipelineError { msg: "rejected".into() })
|
||||
});
|
||||
p.add(reject);
|
||||
assert!(matches!(p.apply("x".to_string()).await, Err(_)));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user