feat: add 4 new crates (queue, tracing, http, cli) + enhancements to existing crates

New crates:
- mytheclipse-queue: unified job queue with WorkerPool, retry/backoff, DLQ
  Backends: in-memory (default), Redis, NATS, PostgreSQL
- mytheclipse-tracing: tracing subscriber layers with env filter + OTLP/Jaeger export
- mytheclipse-http: HTTP client/server with timeout + tracing, axum server
- mytheclipse-cli: CLI framework with clap derive, built-in subcommands

Enhancements to existing crates:
- mytheclipse-core: SemaphorePool, LeaderElection, HealthRegistry/HealthCheck
- mytheclipse-cache: AutoRefreshCache (bg refresh on miss), CacheMetrics
- mytheclipse-config: ConfigSchema for JSON Schema generation
- mytheclipse-storage: MultipartUploadDriver trait
- mytheclipse-crypto: PASETO v4.local token support

All features compile with --all-features; tests + clippy pass clean.
This commit is contained in:
asepharyana
2026-08-29 16:05:02 +07:00
parent 2596b357ec
commit 8106f8943e
53 changed files with 3588 additions and 44 deletions
+53
View File
@@ -0,0 +1,53 @@
//! Errors returned by queue and job operations.
/// Errors from queue-level operations (enqueue, dequeue, etc.).
#[derive(Debug)]
pub enum QueueError {
/// A transport or backend connection error.
Connection(String),
/// The requested topic/queue does not exist or is unavailable.
NotFound(String),
/// A serialization error.
Serialization(String),
/// A timeout occurred while waiting for an operation.
Timeout,
}
impl std::fmt::Display for QueueError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Connection(s) => write!(f, "queue connection error: {s}"),
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"),
}
}
}
impl std::error::Error for QueueError {}
/// Errors from individual job processing.
#[derive(Debug)]
pub enum JobError {
/// The job could not be acknowledged.
AckFailed(String),
/// The job could not be moved to the dead-letter queue.
DlqFailed(String),
/// The job exceeded its maximum retry count.
MaxRetriesExceeded,
/// The job payload could not be decoded.
InvalidPayload(String),
}
impl std::fmt::Display for JobError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::AckFailed(s) => write!(f, "ack failed: {s}"),
Self::DlqFailed(s) => write!(f, "dlq move failed: {s}"),
Self::MaxRetriesExceeded => write!(f, "max retries exceeded"),
Self::InvalidPayload(s) => write!(f, "invalid payload: {s}"),
}
}
}
impl std::error::Error for JobError {}
+162
View File
@@ -0,0 +1,162 @@
//! In-process job queue using `tokio::sync::Mutex` + `Notify`.
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use tokio::sync::{Mutex, Notify};
use crate::error::{JobError, QueueError};
use crate::job::{Job, JobId};
use crate::traits::Queue;
struct TopicQueue {
jobs: Mutex<Vec<Job>>,
notify: Notify,
}
impl std::fmt::Debug for TopicQueue {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TopicQueue")
.field("jobs_len", &self.jobs.try_lock().map(|j| j.len()).unwrap_or(0))
.finish()
}
}
struct Inner {
topics: std::collections::HashMap<String, Arc<TopicQueue>>,
}
impl std::fmt::Debug for Inner {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let keys: Vec<&String> = self.topics.keys().collect();
f.debug_struct("Inner").field("topics", &keys).finish()
}
}
/// An in-memory queue. Each topic is a shared mutex-protected vector + Notify.
#[derive(Clone)]
pub struct InMemoryQueue {
inner: Arc<Mutex<Inner>>,
}
impl std::fmt::Debug for InMemoryQueue {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("InMemoryQueue").finish()
}
}
impl Default for InMemoryQueue {
fn default() -> Self {
Self::new()
}
}
impl InMemoryQueue {
pub fn new() -> Self {
Self {
inner: Arc::new(Mutex::new(Inner {
topics: std::collections::HashMap::new(),
})),
}
}
async fn get_topic(&self, topic: &str) -> Arc<TopicQueue> {
let mut inner = self.inner.lock().await;
inner
.topics
.entry(topic.to_string())
.or_insert_with(|| {
Arc::new(TopicQueue {
jobs: Mutex::new(Vec::new()),
notify: Notify::new(),
})
})
.clone()
}
}
#[async_trait]
impl Queue for InMemoryQueue {
async fn enqueue(&self, topic: &str, payload: Vec<u8>) -> Result<(), QueueError> {
let tq = self.get_topic(topic).await;
tq.jobs.lock().await.push(Job::new(JobId::generate(), topic, payload));
tq.notify.notify_one();
Ok(())
}
async fn dequeue(&self, topic: &str, timeout: Duration) -> Result<Option<Job>, QueueError> {
let tq = self.get_topic(topic).await;
let tq2 = Arc::clone(&tq);
loop {
if let Some(job) = tq.jobs.lock().await.pop() {
return Ok(Some(job));
}
tokio::select! {
_ = tq2.notify.notified() => {}
_ = tokio::time::sleep(timeout) => {
if let Some(job) = tq.jobs.lock().await.pop() {
return Ok(Some(job));
}
return Ok(None);
}
}
}
}
async fn ack(&self, _job: &Job) -> Result<(), JobError> {
Ok(())
}
async fn nack(&self, job: &Job, requeue: bool) -> Result<(), JobError> {
if requeue {
self.enqueue(&job.topic, job.payload.clone())
.await
.map_err(|_| JobError::AckFailed("requeue failed".into()))?;
}
Ok(())
}
async fn dlq_move(&self, topic: &str, job: Job) -> Result<(), QueueError> {
let dlq_topic = format!("dlq:{topic}");
let tq = self.get_topic(&dlq_topic).await;
tq.jobs.lock().await.push(job);
tq.notify.notify_one();
Ok(())
}
async fn len(&self, topic: &str) -> Result<u64, QueueError> {
let tq = self.get_topic(topic).await;
let guard = tq.jobs.lock().await;
Ok(guard.len() as u64)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn enqueue_dequeue_roundtrip() {
let q = InMemoryQueue::new();
q.enqueue("test", b"hello".to_vec()).await.unwrap();
let job = q.dequeue("test", Duration::from_millis(500)).await.unwrap().unwrap();
assert_eq!(job.payload, b"hello");
assert_eq!(job.topic, "test");
}
#[tokio::test]
async fn empty_returns_none() {
let q = InMemoryQueue::new();
let result = q.dequeue("none", Duration::from_millis(50)).await.unwrap();
assert!(result.is_none());
}
#[tokio::test]
async fn len_tracks_jobs() {
let q = InMemoryQueue::new();
q.enqueue("t", vec![1]).await.unwrap();
q.enqueue("t", vec![2]).await.unwrap();
assert_eq!(q.len("t").await.unwrap(), 2);
}
}
+54
View File
@@ -0,0 +1,54 @@
//! The `Job` type and its metadata.
//!
//! A minimal, dependency-free in-process job queue using `tokio::sync::mpsc`.
use std::time::{Duration, SystemTime};
/// A unique identifier for a queued job.
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct JobId(pub String);
impl JobId {
pub fn generate() -> Self {
Self(uuid::Uuid::new_v4().to_string())
}
}
impl std::fmt::Display for JobId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
/// A unit of queued work.
#[derive(Debug, Clone)]
pub struct Job {
/// The unique ID of this job.
pub id: JobId,
/// The topic/queue name the job was delivered from.
pub topic: String,
/// The raw payload bytes.
pub payload: Vec<u8>,
/// How many times this job has been attempted (0 = first attempt).
pub attempt: u32,
/// When the job was first enqueued.
pub enqueued_at: SystemTime,
/// When the job was delivered to the worker (None if not yet delivered).
pub delivered_at: Option<SystemTime>,
/// Optional visibility timeout - after this the job becomes visible again.
pub visibility_timeout: Option<Duration>,
}
impl Job {
pub fn new(id: JobId, topic: &str, payload: Vec<u8>) -> Self {
Self {
id,
topic: topic.to_string(),
payload,
attempt: 0,
enqueued_at: SystemTime::now(),
delivered_at: Some(SystemTime::now()),
visibility_timeout: None,
}
}
}
+63
View File
@@ -0,0 +1,63 @@
//! # mytheclipse-queue
//!
//! A unified job queue abstraction so your background work isn't locked to one
//! transport. Provides a single `Queue` trait, `Job` type, and `WorkerPool`
//! executor with configurable retry/backoff, concurrency, and a dead-letter
//! queue — behind pluggable backends:
//!
//! - **In-memory** (default) — `tokio::sync::mpsc` + task spawning, no external service.
//! - **Redis** (`redis`) — LIST-based queue with atomic moves.
//! - **NATS JetStream** (`nats`) — durable consumer with ACK/NACK.
//! - **PostgreSQL** (`postgres`) — `SKIP LOCKED` polling.
//!
//! ## Quick Start
//!
//! ```toml
//! [dependencies]
//! mytheclipse-queue = "0.2"
//! ```
//!
//! ```ignore
//! use mytheclipse_queue::{InMemoryQueue, WorkerPool, JobHandler, Job};
//! ...
//! let queue = InMemoryQueue::new();
//! queue.enqueue("email", b"hello".to_vec()).await?;
//!
//! fn make_handler() -> impl JobHandler {
//! struct PrintHandler;
//! impl JobHandler for PrintHandler {
//! fn handle(&self, job: Job) -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<(), mytheclipse_queue::JobError>> + Send>> {
//! Box::pin(async move {
//! println!("payload: {:?}", job.payload);
//! Ok(())
//! })
//! }
//! }
//! PrintHandler
//! }
//!
//! # #[tokio::main]
//! # async fn main() -> Result<(), Box<dyn std::error::Error>> {
//! let queue = InMemoryQueue::new();
//! queue.enqueue("email", b"hello".to_vec()).await?;
//!
//! let pool = WorkerPool::new(queue, 4);
//! pool.start("email", make_handler());
//! # Ok(())
//! # }
//! ```
pub mod traits;
pub mod job;
pub mod worker;
pub mod error;
#[cfg(feature = "in-memory")]
pub mod in_memory;
#[cfg(feature = "in-memory")]
pub use in_memory::InMemoryQueue;
pub use traits::Queue;
pub use job::{Job, JobId};
pub use worker::{WorkerPool, WorkerConfig, JobHandler, JobFuture};
pub use error::{QueueError, JobError};
+50
View File
@@ -0,0 +1,50 @@
//! The core `Queue` trait and supporting types.
use async_trait::async_trait;
use std::time::Duration;
use crate::job::Job;
use crate::error::{QueueError, JobError};
/// A handle to a single unit of queued work.
///
/// `Job` carries the raw payload (arbitrary bytes — caller decides encoding)
/// plus metadata the queue implementation fills in (ID, enqueue time, retry
/// count). `ack`/`nack` are only valid on backends that support explicit
/// acknowledgment (NATS, Redis BLPOP-with-confirm). For in-memory and Postgres
/// backends, the worker auto-acknowledges on `Ok` and auto-requeues on `Err`.
/// A trait for enqueueing and dequeueing jobs.
///
/// Implementations must be `Send + Sync`. Each backend provides its own factory
/// (e.g. `InMemoryQueue::new()`, `RedisQueue::connect(url)`).
#[async_trait]
pub trait Queue: Send + Sync {
/// Enqueues `payload` onto `topic`.
async fn enqueue(&self, topic: &str, payload: Vec<u8>) -> Result<(), QueueError>;
/// Dequeues the next job from `topic`, waiting up to `timeout`.
///
/// Returns `None` on timeout when the queue is empty and no job arrives
/// within the window.
async fn dequeue(&self, topic: &str, timeout: Duration) -> Result<Option<Job>, QueueError>;
/// Acknowledges a job as successfully processed.
async fn ack(&self, job: &Job) -> Result<(), JobError>;
/// Negative-acknowledges a job. If `requeue` is true the job goes back
/// onto the queue; if false it moves to the dead-letter queue (if
/// configured).
async fn nack(&self, job: &Job, requeue: bool) -> Result<(), JobError>;
/// Moves a job to the dead-letter queue for `topic`.
async fn dlq_move(&self, topic: &str, job: Job) -> Result<(), QueueError>;
/// Number of messages currently waiting in `topic`.
async fn len(&self, topic: &str) -> Result<u64, QueueError>;
/// Whether the queue is empty.
async fn is_empty(&self, topic: &str) -> Result<bool, QueueError> {
Ok(self.len(topic).await? == 0)
}
}
+168
View File
@@ -0,0 +1,168 @@
//! Worker pool for processing queued jobs with retry, backoff, and graceful shutdown.
use std::pin::Pin;
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use tokio::sync::Semaphore;
use crate::error::JobError;
use crate::job::Job;
use crate::traits::Queue;
/// A future returned by a job handler.
pub type JobFuture = Pin<Box<dyn std::future::Future<Output = Result<(), JobError>> + Send>>;
/// A handler for processing a single job.
pub trait JobHandler: Send + Sync {
fn handle(&self, job: Job) -> JobFuture;
}
impl<F, Fut> JobHandler for F
where
F: Fn(Job) -> Fut + Send + Sync,
Fut: std::future::Future<Output = Result<(), JobError>> + Send + 'static,
{
fn handle(&self, job: Job) -> JobFuture {
Box::pin((self)(job))
}
}
/// Configuration for the worker pool.
#[derive(Debug, Clone)]
pub struct WorkerConfig {
/// Maximum concurrent job handlers.
pub concurrency: usize,
/// Maximum number of retry attempts (0 = no retries).
pub max_retries: u32,
/// Base delay for exponential backoff between retries.
pub retry_base_delay: Duration,
/// Maximum delay cap for retry backoff.
pub retry_max_delay: Duration,
/// Backoff multiplier.
pub retry_factor: f64,
/// Visibility timeout for in-progress jobs.
pub visibility_timeout: Duration,
/// Polling interval when a queue is empty.
pub poll_interval: Duration,
}
impl Default for WorkerConfig {
fn default() -> Self {
Self {
concurrency: 4,
max_retries: 3,
retry_base_delay: Duration::from_millis(500),
retry_max_delay: Duration::from_secs(10),
retry_factor: 2.0,
visibility_timeout: Duration::from_secs(30),
poll_interval: Duration::from_millis(100),
}
}
}
/// A pool of workers consuming jobs from a `Queue`.
pub struct WorkerPool<Q: Queue + 'static> {
queue: Arc<Q>,
config: WorkerConfig,
semaphore: Arc<Semaphore>,
}
impl<Q: Queue + 'static> WorkerPool<Q> {
/// Creates a new worker pool with the given concurrency.
pub fn new(queue: Q, concurrency: usize) -> Self {
Self::with_config(queue, WorkerConfig {
concurrency,
..Default::default()
})
}
/// Creates a new worker pool with explicit configuration.
pub fn with_config(queue: Q, config: WorkerConfig) -> Self {
let sem = Arc::new(Semaphore::new(config.concurrency.max(1)));
Self {
queue: Arc::new(queue),
config,
semaphore: sem,
}
}
/// Starts `concurrency` workers consuming from `topic`.
pub fn start<H>(&self, topic: &str, handler: H)
where
H: JobHandler + 'static,
{
let queue = Arc::clone(&self.queue);
let config = self.config.clone();
let semaphore = Arc::clone(&self.semaphore);
let handler: Arc<dyn JobHandler> = Arc::new(handler);
let topic_owned = topic.to_string();
for _ in 0..config.concurrency {
let q = Arc::clone(&queue);
let sem = Arc::clone(&semaphore);
let h = Arc::clone(&handler);
let cfg = config.clone();
let topic_inner = topic_owned.clone();
tokio::spawn(async move {
loop {
match q.dequeue(&topic_inner, cfg.poll_interval).await {
Ok(Some(job)) => {
let _permit = sem.clone().acquire_owned().await;
let q2 = Arc::clone(&q);
let h2 = Arc::clone(&h);
let cfg2 = cfg.clone();
let t2 = topic_inner.clone();
let j2 = job.clone();
tokio::spawn(async move {
let fut = h2.handle(j2.clone());
match fut.await {
Ok(()) => {
let _ = q2.ack(&j2).await;
}
Err(_) => {
if j2.attempt < cfg2.max_retries {
let _ = q2.nack(&j2, true).await;
} else {
let _ = q2.dlq_move(&t2, j2.clone()).await;
let _ = q2.nack(&j2, false).await;
}
}
}
});
}
Ok(None) => {}
Err(e) => {
tracing::error!("queue error: {e}");
tokio::time::sleep(cfg.poll_interval).await;
}
}
}
});
}
}
}
/// Computes the (capped) exponential backoff delay.
pub fn retry_delay(config: &WorkerConfig, attempt: u32) -> Duration {
let exponent = attempt as f64;
let computed = config.retry_base_delay.as_millis() as f64
* config.retry_factor.powf(exponent.max(0.0));
let capped = computed.min(config.retry_max_delay.as_millis() as f64);
Duration::from_millis(capped as u64)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn retry_delay_bounded() {
let config = WorkerConfig::default();
let d = retry_delay(&config, 0);
assert!(d <= config.retry_max_delay);
}
}