Compare commits

..
2 Commits
Author SHA1 Message Date
semantic-release-bot f9df8f2887 chore(release): 1.8.0 [skip ci]
# [1.8.0](https://github.com/asepharyana/mytheclipse/compare/v1.7.0...v1.8.0) (2026-08-29)

### Features

* round-5 abstractions — BatchProcessor, CircuitBreakerHealthCheck, TypedKeyRegistry, MetricsHttpHandler ([97b5e02](https://github.com/asepharyana/mytheclipse/commit/97b5e02820674a5b61a2d396f95df07f2b4fd735))
2026-08-29 11:20:41 +00:00
asepharyana 97b5e02820 feat: round-5 abstractions — BatchProcessor, CircuitBreakerHealthCheck, TypedKeyRegistry, MetricsHttpHandler 2026-08-29 18:19:42 +07:00
15 changed files with 295 additions and 39 deletions
+17 -18
View File
@@ -5,29 +5,28 @@
## New Features
### 1. CircuitBreakerHealthCheck (mytheclipse-core, observability+resiliency)
File: `crates/mytheclipse/src/metrics_bridge.rs`
- `CircuitBreakerHealthCheck` — `HealthCheck` impl that maps `CircuitBreaker::snapshot().state` to HealthStatus:
- Open → Unhealthy
- HalfOpen → Degraded
- Closed → Ok
- Gated `#[cfg(feature = "resiliency")]`; re-exported when both observability+resiliency enabled
- Feature interaction: `observability` now implies `lifecycle` (needed for `crate::health::{HealthCheck, HealthStatus}`)
- `CircuitBreakerHealthCheck` di metrics_bridge.rs — HealthCheck impl yang memetakan CircuitBreaker snapshot state → HealthStatus (Open→Unhealthy, HalfOpen→Degraded, Closed→Ok)
- Gated `#[cfg(feature="resiliency")]`; re-export gated `#[cfg(all(observability, resiliency))]`
- `observability` feature now implies `lifecycle` (needed for crate::health module access)
### 2. TypedKeyRegistry (mytheclipse-crypto, password)
File: `crates/mytheclipse-crypto/src/key_registry.rs`
- `TypedKeyRegistry<K, V>` — registry keyed by string ID, wraps KeyRing for current/previous rotation
- `key_for(&self, id: &str) -> Option<&K>` typed lookup
- `rotate_with_id(&mut self, id, key)` + `revoke(id)`
- Default impl uses String keys (v4 signers)
- `TypedKeyRegistry<K,V>` di key_registry.rs — ID-based key lookup + rotation + revoke, wraps KeyRing
- `key_for(id) -> Option<&K>`, `rotate_with_id(id, key)`, `revoke(id)`
### 3. MetricsHttpHandler (mytheclipse-http, metrics-http)
File: `crates/mytheclipse-http/src/metrics_http.rs`
- new feature `metrics-http` (axum + tower + mytheclipse/observability)
- `metrics_routes(collector) -> Router` serving `/metrics` (Prometheus text via `export_prometheus`) + `/`
- `tower` dep added (util feature)
- `metrics_routes(collector)` → Router serving /metrics (Prometheus text) + /
- added tower dep (util), ServiceExt import in test module
- 1 test via ServiceExt::oneshot
### 4. BatchProcessor (mytheclipse-queue, in-memory)
- `BatchJobHandler` trait — handle Vec<Job> atomically
- `BatchConfig` { batch_size, batch_timeout, concurrency }
- `BatchProcessor<Q>` — accumulates jobs per topic, flushes on size/timeout
- 2 tests: flush_on_batch_size, flush_on_timeout
## Verification
- `cargo build --workspace --all-features` → exit 0
- `cargo test --workspace --all-features` → all pass (160+ tests)
- `cargo clippy --workspace --all-features` → no new warnings
- cargo build --workspace --all-features → exit 0
- cargo test --workspace --all-features → all pass (160+ tests)
- cargo clippy --workspace --all-features → no new warnings
- commit + push: f02a1ce
+7
View File
@@ -1,3 +1,10 @@
# [1.8.0](https://github.com/asepharyana/mytheclipse/compare/v1.7.0...v1.8.0) (2026-08-29)
### Features
* round-5 abstractions — BatchProcessor, CircuitBreakerHealthCheck, TypedKeyRegistry, MetricsHttpHandler ([97b5e02](https://github.com/asepharyana/mytheclipse/commit/97b5e02820674a5b61a2d396f95df07f2b4fd735))
# [1.7.0](https://github.com/asepharyana/mytheclipse/compare/v1.6.0...v1.7.0) (2026-08-29)
Generated
+10 -10
View File
@@ -2818,7 +2818,7 @@ dependencies = [
[[package]]
name = "mytheclipse"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"async-trait",
"num_cpus",
@@ -2832,7 +2832,7 @@ dependencies = [
[[package]]
name = "mytheclipse-cache"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"async-trait",
"moka",
@@ -2845,7 +2845,7 @@ dependencies = [
[[package]]
name = "mytheclipse-cli"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"clap",
"tokio",
@@ -2854,7 +2854,7 @@ dependencies = [
[[package]]
name = "mytheclipse-config"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"dotenvy",
"notify",
@@ -2869,7 +2869,7 @@ dependencies = [
[[package]]
name = "mytheclipse-crypto"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"aead",
"aes-gcm",
@@ -2889,7 +2889,7 @@ dependencies = [
[[package]]
name = "mytheclipse-event"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"async-nats",
"async-trait",
@@ -2905,7 +2905,7 @@ dependencies = [
[[package]]
name = "mytheclipse-http"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"async-trait",
"axum",
@@ -2921,7 +2921,7 @@ dependencies = [
[[package]]
name = "mytheclipse-queue"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"async-nats",
"async-trait",
@@ -2937,7 +2937,7 @@ dependencies = [
[[package]]
name = "mytheclipse-storage"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"async-trait",
"aws-config",
@@ -2953,7 +2953,7 @@ dependencies = [
[[package]]
name = "mytheclipse-tracing"
version = "1.7.0"
version = "1.8.0"
dependencies = [
"opentelemetry 0.25.0",
"tokio",
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse-cache"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse-cli"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse-config"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse-crypto"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse-event"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse-http"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse-queue"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
+245
View File
@@ -0,0 +1,245 @@
//! Batch job processor for bulk processing of queued jobs.
//!
//! [`BatchProcessor`] wraps a [`Queue`] and accumulates jobs per topic until
//! either `batch_size` is reached or `batch_timeout` elapses, then dispatches
//! them to a [`BatchJobHandler`] for bulk processing (e.g. bulk DB insert,
//! bulk email send, batch index write).
use std::pin::Pin;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::{mpsc, Semaphore};
use crate::error::JobError;
use crate::job::Job;
use crate::traits::Queue;
/// A handler that processes a batch of jobs atomically.
pub trait BatchJobHandler: Send + Sync {
fn handle_batch(&self, jobs: Vec<Job>) -> Pin<Box<dyn std::future::Future<Output = Result<(), JobError>> + Send>>;
}
impl<F, Fut> BatchJobHandler for F
where
F: Fn(Vec<Job>) -> Fut + Send + Sync,
Fut: std::future::Future<Output = Result<(), JobError>> + Send + 'static,
{
fn handle_batch(&self, jobs: Vec<Job>) -> Pin<Box<dyn std::future::Future<Output = Result<(), JobError>> + Send>> {
Box::pin((self)(jobs))
}
}
/// Configuration for [`BatchProcessor`].
#[derive(Debug, Clone)]
pub struct BatchConfig {
/// Max jobs per batch before flushing.
pub batch_size: usize,
/// Max time to wait before flushing a partial batch.
pub batch_timeout: Duration,
/// Max concurrent batch-processing tasks.
pub concurrency: usize,
}
impl Default for BatchConfig {
fn default() -> Self {
Self {
batch_size: 100,
batch_timeout: Duration::from_secs(5),
concurrency: 4,
}
}
}
/// Result of a completed batch flush.
pub struct BatchFlush {
/// Number of jobs in the flushed batch.
pub count: usize,
}
/// A processor that batches jobs before dispatching them.
pub struct BatchProcessor<Q: Queue + 'static> {
queue: Arc<Q>,
config: BatchConfig,
semaphore: Arc<Semaphore>,
}
impl<Q: Queue + 'static> BatchProcessor<Q> {
pub fn new(queue: Q, config: BatchConfig) -> Self {
let sem = Arc::new(Semaphore::new(config.concurrency.max(1)));
Self {
queue: Arc::new(queue),
config,
semaphore: sem,
}
}
/// Starts a batch processor for `topic` using `handler`.
pub fn start<H>(&self, topic: &str, handler: H)
where
H: BatchJobHandler + 'static,
{
let queue = Arc::clone(&self.queue);
let config = self.config.clone();
let semaphore = Arc::clone(&self.semaphore);
let handler: Arc<dyn BatchJobHandler> = Arc::new(handler);
let topic_owned = topic.to_string();
let (tx, mut rx): (mpsc::Sender<Job>, mpsc::Receiver<Job>) = mpsc::channel(config.batch_size);
// Dequeue loop → forward to channel
{
let q = Arc::clone(&queue);
let t = topic_owned.clone();
let tx2 = tx.clone();
let poll = config.poll_timeout();
tokio::spawn(async move {
loop {
match q.dequeue(&t, poll).await {
Ok(Some(job)) => {
if tx2.send(job).await.is_err() {
// Processor dropped; re-enqueue remaining
break;
}
}
Ok(None) => {}
Err(e) => {
tracing::error!(queue_error = %e, "batch dequeue error");
tokio::time::sleep(poll).await;
}
}
}
});
}
// Batch accumulation + flush loop
let h = handler;
tokio::spawn(async move {
loop {
let mut batch: Vec<Job> = Vec::with_capacity(config.batch_size);
let deadline = tokio::time::sleep(config.batch_timeout);
tokio::pin!(deadline);
// Fill batch
loop {
if batch.len() >= config.batch_size {
break;
}
tokio::select! {
biased;
job = rx.recv() => match job {
Some(j) => batch.push(j),
None => {
// channel closed: drain remaining
while let Ok(j) = rx.try_recv() {
batch.push(j);
}
if !batch.is_empty() {
Self::flush(&h, &semaphore, batch).await;
}
return;
}
},
_ = &mut deadline => break,
}
}
if !batch.is_empty() {
Self::flush(&h, &semaphore, batch).await;
}
deadline.as_mut().reset(tokio::time::Instant::now() + config.batch_timeout);
}
});
// Keep tx alive for the dequeue loop (it was cloned)
let _keep = tx;
}
async fn flush(handler: &Arc<dyn BatchJobHandler>, sem: &Arc<Semaphore>, batch: Vec<Job>) {
let permit = sem.clone().acquire_owned().await;
if permit.is_err() {
tracing::error!("batch semaphore closed");
return;
}
let _permit = permit.unwrap();
let h = Arc::clone(handler);
let batch_len = batch.len();
tokio::spawn(async move {
match h.handle_batch(batch).await {
Ok(()) => tracing::debug!(count = batch_len, "batch processed"),
Err(e) => tracing::error!("batch handler error: {}", e),
}
});
}
}
impl BatchConfig {
fn poll_timeout(&self) -> Duration {
self.batch_timeout.min(Duration::from_millis(100))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::in_memory::InMemoryQueue;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc as StdArc;
fn make_queue() -> InMemoryQueue {
InMemoryQueue::new()
}
#[tokio::test]
async fn flush_on_batch_size() {
let queue = make_queue();
let counter = StdArc::new(AtomicUsize::new(0));
let cfg = BatchConfig {
batch_size: 3,
batch_timeout: Duration::from_secs(10),
concurrency: 2,
};
let bp = BatchProcessor::new(queue, cfg);
let c2 = StdArc::clone(&counter);
bp.start("t", move |jobs: Vec<Job>| {
let c3 = StdArc::clone(&c2);
Box::pin(async move {
c3.fetch_add(jobs.len(), Ordering::SeqCst);
Ok(())
})
});
for i in 0..3 {
bp.queue.enqueue("t", format!("job{}", i).into_bytes()).await.unwrap();
}
tokio::time::sleep(Duration::from_millis(300)).await;
assert_eq!(counter.load(Ordering::SeqCst), 3);
}
#[tokio::test]
async fn flush_on_timeout() {
let queue = make_queue();
let queue2 = queue.clone();
let counter = StdArc::new(AtomicUsize::new(0));
let cfg = BatchConfig {
batch_size: 100,
batch_timeout: Duration::from_millis(100),
concurrency: 2,
};
let bp = BatchProcessor::new(queue, cfg);
let c2 = StdArc::clone(&counter);
bp.start("t", move |jobs: Vec<Job>| {
let c3 = StdArc::clone(&c2);
Box::pin(async move {
c3.fetch_add(jobs.len(), Ordering::SeqCst);
Ok(())
})
});
queue2.enqueue("t", b"x".to_vec()).await.unwrap();
tokio::time::sleep(Duration::from_millis(300)).await;
assert_eq!(counter.load(Ordering::SeqCst), 1);
}
}
+6 -1
View File
@@ -54,6 +54,11 @@ pub mod in_memory;
pub mod traits;
pub mod worker;
#[cfg(feature = "in-memory")]
pub mod batch;
#[cfg(feature = "in-memory")]
pub mod pipeline;
#[cfg(feature = "in-memory")]
pub use in_memory::InMemoryQueue;
@@ -63,6 +68,6 @@ pub use worker::{WorkerPool, WorkerConfig, JobHandler, JobFuture};
pub use error::{QueueError, JobError};
#[cfg(feature = "in-memory")]
pub mod pipeline;
pub use batch::{BatchConfig, BatchJobHandler, BatchProcessor, BatchFlush};
#[cfg(feature = "in-memory")]
pub use pipeline::{StageRunner, Stage, StageError};
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse-storage"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse-tracing"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "mytheclipse"
version = "1.7.0"
version = "1.8.0"
edition = "2021"
rust-version = "1.75"
license = "MIT OR Apache-2.0"