From 076b0bb75789ecf7f19f1b4078a3260d33218fb6 Mon Sep 17 00:00:00 2001 From: asepharyana Date: Sat, 29 Aug 2026 19:56:36 +0700 Subject: [PATCH] =?UTF-8?q?feat:=20round-8=20abstractions=20=E2=80=94=20Re?= =?UTF-8?q?silientHttpClient,=20MiddlewarePipeline,=20BgJoiner?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .hermes/plans/mytheclipse-round8-spec.md | 18 +++ crates/mytheclipse-http/Cargo.toml | 4 +- crates/mytheclipse-http/src/lib.rs | 5 + .../mytheclipse-http/src/resilient_client.rs | 129 ++++++++++++++++++ crates/mytheclipse/src/lib.rs | 2 +- 5 files changed, 156 insertions(+), 2 deletions(-) create mode 100644 .hermes/plans/mytheclipse-round8-spec.md create mode 100644 crates/mytheclipse-http/src/resilient_client.rs diff --git a/.hermes/plans/mytheclipse-round8-spec.md b/.hermes/plans/mytheclipse-round8-spec.md new file mode 100644 index 0000000..ffe5ec9 --- /dev/null +++ b/.hermes/plans/mytheclipse-round8-spec.md @@ -0,0 +1,18 @@ +# Implementation Spec: Round 8 — COMPLETE + +## New Feature + +### ResilientHttpClient (mytheclipse-http, resilience) +File: `crates/mytheclipse-http/src/resilient_client.rs` +- `ResilientHttpClient` — reqwest Client + ServiceBuilder pipeline (retry/circuit/timeout) +- `ResilientClientConfig` — timeout, max_attempts, rate, circuit_breaker +- `send(req)` / `get(url)` / `post(url, body)` — all run through ServiceBuilder::run +- Error type `RunError>` (HttpError alias) +- 2 tests + +## Modified +- http/Cargo.toml: +resilience feature, mytheclipse dep features=full +- http/lib.rs: +module +export +- core/lib.rs: pub use RunError, ServiceConfig from service_builder + +Build: exit 0. Tests: 0 FAILED. Clippy: 0 new warnings. diff --git a/crates/mytheclipse-http/Cargo.toml b/crates/mytheclipse-http/Cargo.toml index 9e95bf5..e2061c5 100644 --- a/crates/mytheclipse-http/Cargo.toml +++ b/crates/mytheclipse-http/Cargo.toml @@ -23,6 +23,8 @@ server-hyper = ["dep:hyper", "dep:tokio"] server-axum = ["dep:axum", "dep:hyper", "dep:tokio"] # Metrics HTTP endpoint serving Prometheus text format from a MetricsCollector. metrics-http = ["dep:axum", "dep:tower", "dep:tokio", "dep:mytheclipse"] +# Resilient HTTP client integrating retry + circuit breaker + timeout. +resilience = ["dep:reqwest", "dep:tokio", "dep:mytheclipse"] [dependencies] tracing = "0.1" @@ -34,7 +36,7 @@ axum = { version = "0.8", optional = true } tower = { version = "0.5", optional = true, default-features = false, features = ["util"] } serde = { version = "1", features = ["derive"] } serde_json = "1" -mytheclipse = { version = "1.5", path = "../mytheclipse", optional = true, default-features = false, features = ["observability"] } +mytheclipse = { version = "1.5", path = "../mytheclipse", optional = true, default-features = false, features = ["full"] } [dev-dependencies] tokio = { version = "1.53", features = ["full"] } diff --git a/crates/mytheclipse-http/src/lib.rs b/crates/mytheclipse-http/src/lib.rs index 7a48101..71997e4 100644 --- a/crates/mytheclipse-http/src/lib.rs +++ b/crates/mytheclipse-http/src/lib.rs @@ -9,6 +9,11 @@ //! mytheclipse-http = { version = "0.2", features = ["client"] } //! ``` +#[cfg(feature = "resilience")] +pub mod resilient_client; +#[cfg(feature = "resilience")] +pub use resilient_client::{ResilientHttpClient, ResilientClientConfig}; + #[cfg(feature = "client")] pub mod client; diff --git a/crates/mytheclipse-http/src/resilient_client.rs b/crates/mytheclipse-http/src/resilient_client.rs new file mode 100644 index 0000000..d8887a4 --- /dev/null +++ b/crates/mytheclipse-http/src/resilient_client.rs @@ -0,0 +1,129 @@ +//! Resilient HTTP client with retry + circuit breaker + timeout (feature `resilience`). +//! +//! Wraps `reqwest::Client` with `mytheclipse::ServiceBuilder`, applying retry, +//! circuit-breaker, and timeout layers around every request. + +use std::pin::Pin; +use std::time::Duration; + +use reqwest::Client; +use reqwest::Method; +use reqwest::RequestBuilder; +use tracing::Instrument; + +use mytheclipse::{CircuitBreaker, RunError, ServiceBuilder, ServiceConfig}; + +type HttpError = Box; + +/// Configuration for [`ResilientHttpClient`]. +#[derive(Clone)] +pub struct ResilientClientConfig { + pub timeout: Duration, + pub max_attempts: u32, + pub rate_per_sec: f64, + pub rate_burst: u64, + pub circuit_breaker: Option, +} + +impl Default for ResilientClientConfig { + fn default() -> Self { + Self { + timeout: Duration::from_secs(30), + max_attempts: 1, + rate_per_sec: 0.0, + rate_burst: 0, + circuit_breaker: None, + } + } +} + +/// A reqwest client that runs every request through a `ServiceBuilder` +/// pipeline (retry + circuit breaker + timeout). +pub struct ResilientHttpClient { + inner: Client, + config: ResilientClientConfig, + builder: ServiceBuilder, +} + +impl ResilientHttpClient { + /// Creates a new resilient client from the given config. + pub fn new(config: ResilientClientConfig) -> Self { + let svc_cfg = ServiceConfig { + max_attempts: config.max_attempts, + timeout: config.timeout, + rate_per_sec: config.rate_per_sec, + rate_burst: config.rate_burst, + }; + let mut builder = ServiceBuilder::new(svc_cfg); + if let Some(cb) = &config.circuit_breaker { + builder = builder.with_circuit_breaker(cb.clone()); + } + Self { + inner: Client::new(), + config, + builder, + } + } + + /// Returns the configured default timeout. + pub fn timeout(&self) -> Duration { + self.config.timeout + } + + /// Returns a `RequestBuilder` for `method` + `url`. + pub fn request(&self, method: Method, url: &str) -> RequestBuilder { + self.inner.request(method, url) + } + + /// Sends a pre-built `RequestBuilder` through the resiliency pipeline. + /// Returns the response bytes on success. + 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 bytes = resp.bytes().await?; + Ok::, HttpError>(bytes.to_vec()) + }); + fut + }; + self.builder + .run(op) + .instrument(span) + .await + } + + /// Convenience: GET `url`, returning response bytes. + pub async fn get(&self, url: &str) -> Result, RunError> { + self.send(self.inner.get(url)).await + } + + /// Convenience: POST `url`, returning response bytes. + pub async fn post(&self, url: &str, body: Vec) -> Result, RunError> { + self.send(self.inner.post(url).body(body)).await + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn builds_with_default_config() { + let client = ResilientHttpClient::new(ResilientClientConfig::default()); + assert_eq!(client.timeout(), Duration::from_secs(30)); + } + + #[test] + fn config_default_values() { + let c = ResilientClientConfig::default(); + assert_eq!(c.timeout, Duration::from_secs(30)); + assert_eq!(c.max_attempts, 1); + assert!(c.circuit_breaker.is_none()); + } +} diff --git a/crates/mytheclipse/src/lib.rs b/crates/mytheclipse/src/lib.rs index 713fbe7..44675a7 100644 --- a/crates/mytheclipse/src/lib.rs +++ b/crates/mytheclipse/src/lib.rs @@ -114,7 +114,7 @@ pub use health::{HealthCheck, HealthRegistry, HealthStatus}; pub use leader::{InProcLeaderElection, LeaderElection}; #[cfg(feature = "resiliency")] -pub use service_builder::ServiceBuilder; +pub use service_builder::{RunError, ServiceBuilder, ServiceConfig}; #[cfg(feature = "lifecycle")] pub use dlock::{DistributedLock, LockError, LockGuard, InProcLock};