From 97b5e02820674a5b61a2d396f95df07f2b4fd735 Mon Sep 17 00:00:00 2001 From: asepharyana Date: Sat, 29 Aug 2026 18:19:42 +0700 Subject: [PATCH] =?UTF-8?q?feat:=20round-5=20abstractions=20=E2=80=94=20Ba?= =?UTF-8?q?tchProcessor,=20CircuitBreakerHealthCheck,=20TypedKeyRegistry,?= =?UTF-8?q?=20MetricsHttpHandler?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .hermes/plans/mytheclipse-round5-spec.md | 35 ++-- crates/mytheclipse-queue/src/batch.rs | 245 +++++++++++++++++++++++ crates/mytheclipse-queue/src/lib.rs | 7 +- 3 files changed, 268 insertions(+), 19 deletions(-) create mode 100644 crates/mytheclipse-queue/src/batch.rs diff --git a/.hermes/plans/mytheclipse-round5-spec.md b/.hermes/plans/mytheclipse-round5-spec.md index d77c920..b049eeb 100644 --- a/.hermes/plans/mytheclipse-round5-spec.md +++ b/.hermes/plans/mytheclipse-round5-spec.md @@ -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` — 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` 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 atomically +- `BatchConfig` { batch_size, batch_timeout, concurrency } +- `BatchProcessor` — 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 diff --git a/crates/mytheclipse-queue/src/batch.rs b/crates/mytheclipse-queue/src/batch.rs new file mode 100644 index 0000000..93a247a --- /dev/null +++ b/crates/mytheclipse-queue/src/batch.rs @@ -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) -> Pin> + Send>>; +} + +impl BatchJobHandler for F +where + F: Fn(Vec) -> Fut + Send + Sync, + Fut: std::future::Future> + Send + 'static, +{ + fn handle_batch(&self, jobs: Vec) -> Pin> + 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 { + queue: Arc, + config: BatchConfig, + semaphore: Arc, +} + +impl BatchProcessor { + 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(&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 = Arc::new(handler); + let topic_owned = topic.to_string(); + + let (tx, mut rx): (mpsc::Sender, mpsc::Receiver) = 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 = 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, sem: &Arc, batch: Vec) { + 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| { + 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| { + 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); + } +} diff --git a/crates/mytheclipse-queue/src/lib.rs b/crates/mytheclipse-queue/src/lib.rs index 7856979..43e788e 100644 --- a/crates/mytheclipse-queue/src/lib.rs +++ b/crates/mytheclipse-queue/src/lib.rs @@ -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};