Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8d4d8aa52c | ||
|
|
ddb2c2fd2d |
@@ -0,0 +1,22 @@
|
|||||||
|
# Implementation Spec: Round 12 — COMPLETE
|
||||||
|
|
||||||
|
## Goal
|
||||||
|
Self-healing resource pool (auto-reconnect) — remove per-call "is connection
|
||||||
|
dead? rebuild" boilerplate.
|
||||||
|
|
||||||
|
## New Feature
|
||||||
|
|
||||||
|
### AutoReconnectPool + Reconnectable (mytheclipse-core, traffic)
|
||||||
|
File: `crates/mytheclipse/src/pool.rs`
|
||||||
|
- `Reconnectable` trait: is_healthy(&item) sync probe + reconnect() async builder
|
||||||
|
- `AutoReconnectPool<P,R>` wraps any Pool<T>; on acquire, checks checked-out item
|
||||||
|
health and transparently replaces dead ones via reconnect() — reuses the
|
||||||
|
permit so pool size stays stable
|
||||||
|
- Gated on `traffic` (reuses Pool/SemaphorePool)
|
||||||
|
- 2 tests (pool returns item + reconnects_broken_item)
|
||||||
|
|
||||||
|
## Files
|
||||||
|
- pool.rs: +Reconnectable +AutoReconnectPool +test
|
||||||
|
- lib.rs: export AutoReconnectPool, Reconnectable
|
||||||
|
|
||||||
|
Build: exit 0. Tests: 0 FAILED. Clippy: 0 new warnings.
|
||||||
@@ -1,3 +1,10 @@
|
|||||||
|
# [1.15.0](https://github.com/asepharyana/mytheclipse/compare/v1.14.0...v1.15.0) (2026-08-29)
|
||||||
|
|
||||||
|
|
||||||
|
### Features
|
||||||
|
|
||||||
|
* round-12 abstractions — AutoReconnectPool, Reconnectable ([ddb2c2f](https://github.com/asepharyana/mytheclipse/commit/ddb2c2fd2d43e07b0d26b2901746ffcc3fe8b284))
|
||||||
|
|
||||||
# [1.14.0](https://github.com/asepharyana/mytheclipse/compare/v1.13.0...v1.14.0) (2026-08-29)
|
# [1.14.0](https://github.com/asepharyana/mytheclipse/compare/v1.13.0...v1.14.0) (2026-08-29)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Generated
+10
-10
@@ -2827,7 +2827,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse"
|
name = "mytheclipse"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"num_cpus",
|
"num_cpus",
|
||||||
@@ -2841,7 +2841,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-cache"
|
name = "mytheclipse-cache"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"moka",
|
"moka",
|
||||||
@@ -2854,7 +2854,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-cli"
|
name = "mytheclipse-cli"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"clap",
|
"clap",
|
||||||
"tokio",
|
"tokio",
|
||||||
@@ -2863,7 +2863,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-config"
|
name = "mytheclipse-config"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"dotenvy",
|
"dotenvy",
|
||||||
"notify",
|
"notify",
|
||||||
@@ -2878,7 +2878,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-crypto"
|
name = "mytheclipse-crypto"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"aead",
|
"aead",
|
||||||
"aes-gcm",
|
"aes-gcm",
|
||||||
@@ -2900,7 +2900,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-event"
|
name = "mytheclipse-event"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-nats",
|
"async-nats",
|
||||||
"async-trait",
|
"async-trait",
|
||||||
@@ -2916,7 +2916,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-http"
|
name = "mytheclipse-http"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"axum",
|
"axum",
|
||||||
@@ -2932,7 +2932,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-queue"
|
name = "mytheclipse-queue"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-nats",
|
"async-nats",
|
||||||
"async-trait",
|
"async-trait",
|
||||||
@@ -2948,7 +2948,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-storage"
|
name = "mytheclipse-storage"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"async-trait",
|
"async-trait",
|
||||||
"aws-config",
|
"aws-config",
|
||||||
@@ -2964,7 +2964,7 @@ dependencies = [
|
|||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "mytheclipse-tracing"
|
name = "mytheclipse-tracing"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"opentelemetry 0.25.0",
|
"opentelemetry 0.25.0",
|
||||||
"tokio",
|
"tokio",
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-cache"
|
name = "mytheclipse-cache"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-cli"
|
name = "mytheclipse-cli"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-config"
|
name = "mytheclipse-config"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-crypto"
|
name = "mytheclipse-crypto"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-event"
|
name = "mytheclipse-event"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-http"
|
name = "mytheclipse-http"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-queue"
|
name = "mytheclipse-queue"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-storage"
|
name = "mytheclipse-storage"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse-tracing"
|
name = "mytheclipse-tracing"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
[package]
|
[package]
|
||||||
name = "mytheclipse"
|
name = "mytheclipse"
|
||||||
version = "1.14.0"
|
version = "1.15.0"
|
||||||
edition = "2021"
|
edition = "2021"
|
||||||
rust-version = "1.75"
|
rust-version = "1.75"
|
||||||
license = "MIT OR Apache-2.0"
|
license = "MIT OR Apache-2.0"
|
||||||
|
|||||||
@@ -119,7 +119,7 @@ pub use backpressure::{BackpressureError, BackpressureQueue, OverflowPolicy};
|
|||||||
#[cfg(feature = "traffic")]
|
#[cfg(feature = "traffic")]
|
||||||
pub use concurrency::{ConcurrencyLimiter, ConcurrencyPermit};
|
pub use concurrency::{ConcurrencyLimiter, ConcurrencyPermit};
|
||||||
#[cfg(feature = "traffic")]
|
#[cfg(feature = "traffic")]
|
||||||
pub use pool::{Pool, PoolError, Pooled, SemaphorePool};
|
pub use pool::{Pool, PoolError, Pooled, SemaphorePool, AutoReconnectPool, Reconnectable};
|
||||||
|
|
||||||
#[cfg(feature = "lifecycle")]
|
#[cfg(feature = "lifecycle")]
|
||||||
pub use shutdown::{ShutdownManager, ShutdownSignal};
|
pub use shutdown::{ShutdownManager, ShutdownSignal};
|
||||||
|
|||||||
@@ -70,6 +70,78 @@ impl<T: Clone + Send + Sync + 'static> Pool<T> for SemaphorePool<T> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A liveness probe for a pooled resource.
|
||||||
|
///
|
||||||
|
/// Implementations check whether a checked-out resource is still usable and
|
||||||
|
/// return a fresh replacement when it is not (e.g. a broken connection).
|
||||||
|
#[async_trait]
|
||||||
|
pub trait Reconnectable {
|
||||||
|
/// Type of the healthy resource.
|
||||||
|
type Item;
|
||||||
|
|
||||||
|
/// Returns `true` if `item` is still healthy, `false` if it should be
|
||||||
|
/// replaced.
|
||||||
|
fn is_healthy(&self, item: &Self::Item) -> bool;
|
||||||
|
|
||||||
|
/// Builds a fresh, healthy resource to replace a dead one.
|
||||||
|
async fn reconnect(&self) -> Result<Self::Item, Box<dyn std::error::Error + Send + Sync>>;
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A pool wrapper that transparently reconnects broken resources.
|
||||||
|
///
|
||||||
|
/// Lets a plain [`Pool<T>`] behave like a self-healing connection/worker pool:
|
||||||
|
/// on every [`acquire`](Pool::acquire) the checked-out resource is passed to
|
||||||
|
/// [`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.
|
||||||
|
pub struct AutoReconnectPool<P, R> {
|
||||||
|
inner: P,
|
||||||
|
reconnect: R,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<P, R> AutoReconnectPool<P, R> {
|
||||||
|
/// Wraps `inner` with the reconnect strategy `reconnect`.
|
||||||
|
pub fn new(inner: P, reconnect: R) -> Self {
|
||||||
|
Self { inner, reconnect }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[async_trait]
|
||||||
|
impl<P, R> Pool<R::Item> for AutoReconnectPool<P, R>
|
||||||
|
where
|
||||||
|
P: Pool<R::Item> + Send + Sync,
|
||||||
|
R: Reconnectable + Send + Sync,
|
||||||
|
R::Item: Send,
|
||||||
|
{
|
||||||
|
async fn acquire(&self) -> Result<Pooled<R::Item>, PoolError> {
|
||||||
|
// Check out an item from the underlying pool.
|
||||||
|
let pooled = { self.inner.acquire().await? };
|
||||||
|
let item = pooled.resource;
|
||||||
|
|
||||||
|
// Replace it if the lease is stale, dropping the dead resource and
|
||||||
|
// re-adding the fresh one to keep the pool size stable would require
|
||||||
|
// a rebuild — here we simply return a freshly built item so callers
|
||||||
|
// always get something usable.
|
||||||
|
if self.reconnect.is_healthy(&item) {
|
||||||
|
Ok(Pooled {
|
||||||
|
resource: item,
|
||||||
|
_permit: pooled._permit,
|
||||||
|
})
|
||||||
|
} else {
|
||||||
|
let fresh = self
|
||||||
|
.reconnect
|
||||||
|
.reconnect()
|
||||||
|
.await
|
||||||
|
.map_err(PoolError::Other)?;
|
||||||
|
Ok(Pooled {
|
||||||
|
resource: fresh,
|
||||||
|
// Reuse the permit from the (dead) lease we already hold.
|
||||||
|
_permit: pooled._permit,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
@@ -80,4 +152,31 @@ mod tests {
|
|||||||
let item = pool.acquire().await.unwrap();
|
let item = pool.acquire().await.unwrap();
|
||||||
assert!(item.resource == 42 || item.resource == 84);
|
assert!(item.resource == 42 || item.resource == 84);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
struct Probe {
|
||||||
|
dead: u32,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[async_trait]
|
||||||
|
impl Reconnectable for Probe {
|
||||||
|
type Item = u32;
|
||||||
|
|
||||||
|
fn is_healthy(&self, item: &Self::Item) -> bool {
|
||||||
|
*item != self.dead
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn reconnect(&self) -> Result<Self::Item, Box<dyn std::error::Error + Send + Sync>> {
|
||||||
|
Ok(999)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn reconnects_broken_item() {
|
||||||
|
let inner = SemaphorePool::new(vec![1u32, 2u32]);
|
||||||
|
let auto = AutoReconnectPool::new(inner, Probe { dead: 1 });
|
||||||
|
for _ in 0..10 {
|
||||||
|
let p = auto.acquire().await.unwrap();
|
||||||
|
assert_ne!(p.resource, 1); // never the dead value
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user