fix(ci): restore full CI green — test-matrix, clippy, rustfmt, and rustdoc gates
Root cause of the failing CI run was that examples/tests referencing feature-gated items were auto-detected (no required-features), so `--all-targets` compiled them under feature combinations where those modules didn't exist. Fixes: - mytheclipse Cargo.toml: declare the `high_level` example and `race_stress` integration test with required-features = ["full"]; `cargo build --all-targets` now skips them when full is off. This clears the whole test-matrix (workspace default, all-features, and every single-feature config) which all failed on E0432/E0433. - lib.rs: auto_metrics_service depends on service_builder, so regate it behind all(observability, resiliency) instead of observability alone (observability-only build compiled the module without resiliency). - mytheclipse-tracing: gate `pub mod fmt` behind any tracing-subscriber-providing feature so --no-default-features compiles. - mytheclipse-queue: gate `pub mod worker` behind in-memory (worker.rs requires tokio, only provided by in-memory). - clippy -D warnings fixes: deprecated base64 0.22 free fns -> Engine (paseto), unused key field, needless mut (service_builder), unused import/dead var/missing is_empty (bg_join), dead is_expired (dlock), while-let-iterator->for (parallel_map), type_complexity (shutdown_guard), MutexGuard held across await (middleware, now clones Arc'd layers), if-let-Err->is_err (queue), unused CliBuilder fields now wired into clap. - rustdoc -D warnings: resolve retry/MetricsBridge/CircuitBreaker/KeyRing intra-doc links and fix the unparseable lifecycle.rs code fence. - cargo fmt --all to satisfy the Rustfmt gate.
This commit is contained in:
@@ -113,7 +113,7 @@ pub async fn enqueue_with_backpressure<Q: Queue + ?Sized>(
|
||||
) -> Result<usize, BackpressureError> {
|
||||
let mut rejected = 0;
|
||||
for payload in payloads {
|
||||
if let Err(_) = enforcer.try_enqueue(queue, topic, payload).await {
|
||||
if enforcer.try_enqueue(queue, topic, payload).await.is_err() {
|
||||
rejected += 1;
|
||||
}
|
||||
}
|
||||
@@ -142,6 +142,9 @@ mod tests {
|
||||
let _first = reg.global.clone().acquire_owned().await.unwrap();
|
||||
|
||||
let result = reg.try_enqueue(&queue, "t", b"x".to_vec()).await;
|
||||
assert!(matches!(result, Err(BackpressureError::LimitReached { .. })));
|
||||
assert!(matches!(
|
||||
result,
|
||||
Err(BackpressureError::LimitReached { .. })
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,7 +17,10 @@ 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>>;
|
||||
fn handle_batch(
|
||||
&self,
|
||||
jobs: Vec<Job>,
|
||||
) -> Pin<Box<dyn std::future::Future<Output = Result<(), JobError>> + Send>>;
|
||||
}
|
||||
|
||||
impl<F, Fut> BatchJobHandler for F
|
||||
@@ -25,7 +28,10 @@ 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>> {
|
||||
fn handle_batch(
|
||||
&self,
|
||||
jobs: Vec<Job>,
|
||||
) -> Pin<Box<dyn std::future::Future<Output = Result<(), JobError>> + Send>> {
|
||||
Box::pin((self)(jobs))
|
||||
}
|
||||
}
|
||||
@@ -85,7 +91,8 @@ impl<Q: Queue + 'static> BatchProcessor<Q> {
|
||||
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);
|
||||
let (tx, mut rx): (mpsc::Sender<Job>, mpsc::Receiver<Job>) =
|
||||
mpsc::channel(config.batch_size);
|
||||
|
||||
// Dequeue loop → forward to channel
|
||||
{
|
||||
@@ -147,7 +154,9 @@ impl<Q: Queue + 'static> BatchProcessor<Q> {
|
||||
if !batch.is_empty() {
|
||||
Self::flush(&h, &semaphore, batch).await;
|
||||
}
|
||||
deadline.as_mut().reset(tokio::time::Instant::now() + config.batch_timeout);
|
||||
deadline
|
||||
.as_mut()
|
||||
.reset(tokio::time::Instant::now() + config.batch_timeout);
|
||||
}
|
||||
});
|
||||
|
||||
@@ -210,7 +219,10 @@ mod tests {
|
||||
});
|
||||
|
||||
for i in 0..3 {
|
||||
bp.queue.enqueue("t", format!("job{}", i).into_bytes()).await.unwrap();
|
||||
bp.queue
|
||||
.enqueue("t", format!("job{}", i).into_bytes())
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
tokio::time::sleep(Duration::from_millis(300)).await;
|
||||
|
||||
@@ -18,7 +18,10 @@ struct TopicQueue {
|
||||
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))
|
||||
.field(
|
||||
"jobs_len",
|
||||
&self.jobs.try_lock().map(|j| j.len()).unwrap_or(0),
|
||||
)
|
||||
.finish()
|
||||
}
|
||||
}
|
||||
@@ -80,7 +83,10 @@ impl InMemoryQueue {
|
||||
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.jobs
|
||||
.lock()
|
||||
.await
|
||||
.push(Job::new(JobId::generate(), topic, payload));
|
||||
tq.notify.notify_one();
|
||||
Ok(())
|
||||
}
|
||||
@@ -140,7 +146,11 @@ mod tests {
|
||||
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();
|
||||
let job = q
|
||||
.dequeue("test", Duration::from_millis(500))
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(job.payload, b"hello");
|
||||
assert_eq!(job.topic, "test");
|
||||
}
|
||||
|
||||
@@ -11,28 +11,31 @@
|
||||
//! - **PostgreSQL** (`postgres`) — `SKIP LOCKED` polling.
|
||||
|
||||
pub mod error;
|
||||
pub mod job;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod in_memory;
|
||||
pub mod job;
|
||||
pub mod traits;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod worker;
|
||||
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod batch;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod backpressure_enqueue;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod batch;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod rate_limited;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod worker_rate_limited;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub use backpressure_enqueue::{BackpressureEnforcer, BackpressureError, enqueue_with_backpressure};
|
||||
pub use backpressure_enqueue::{
|
||||
enqueue_with_backpressure, BackpressureEnforcer, BackpressureError,
|
||||
};
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub use rate_limited::{RateLimitedQueue, RateLimitQueueError};
|
||||
pub use rate_limited::{RateLimitQueueError, RateLimitedQueue};
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub use worker_rate_limited::RateLimitedWorkerPool;
|
||||
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub mod pipeline;
|
||||
#[cfg(feature = "in-memory")]
|
||||
pub use pipeline::{StageRunner, Stage, StageError};
|
||||
pub use pipeline::{Stage, StageError, StageRunner};
|
||||
|
||||
@@ -75,16 +75,14 @@ where
|
||||
let mut input = input;
|
||||
loop {
|
||||
match input.recv().await {
|
||||
Some(item) => {
|
||||
match stage.process(item).await {
|
||||
Ok(out) => {
|
||||
if output.send(out).await.is_err() {
|
||||
return Err(StageError::ChannelClosed);
|
||||
}
|
||||
Some(item) => match stage.process(item).await {
|
||||
Ok(out) => {
|
||||
if output.send(out).await.is_err() {
|
||||
return Err(StageError::ChannelClosed);
|
||||
}
|
||||
Err(e) => return Err(e),
|
||||
}
|
||||
}
|
||||
Err(e) => return Err(e),
|
||||
},
|
||||
None => return Ok(()),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -111,7 +111,11 @@ impl<Q: Queue + ?Sized> Queue for RateLimitedQueue<Q> {
|
||||
self.inner.enqueue(topic, payload).await
|
||||
}
|
||||
|
||||
async fn dequeue(&self, topic: &str, timeout: Duration) -> Result<Option<crate::job::Job>, QueueError> {
|
||||
async fn dequeue(
|
||||
&self,
|
||||
topic: &str,
|
||||
timeout: Duration,
|
||||
) -> Result<Option<crate::job::Job>, QueueError> {
|
||||
self.inner.dequeue(topic, timeout).await
|
||||
}
|
||||
|
||||
@@ -119,7 +123,11 @@ impl<Q: Queue + ?Sized> Queue for RateLimitedQueue<Q> {
|
||||
self.inner.ack(job).await
|
||||
}
|
||||
|
||||
async fn nack(&self, job: &crate::job::Job, requeue: bool) -> Result<(), crate::error::JobError> {
|
||||
async fn nack(
|
||||
&self,
|
||||
job: &crate::job::Job,
|
||||
requeue: bool,
|
||||
) -> Result<(), crate::error::JobError> {
|
||||
self.inner.nack(job, requeue).await
|
||||
}
|
||||
|
||||
@@ -148,7 +156,7 @@ mod tests {
|
||||
async fn rejects_when_bucket_empty() {
|
||||
let inner = InMemoryQueue::new();
|
||||
let rl = RateLimitedQueue::new(inner, 0.0, 1); // 0 tokens/sec, 1 burst
|
||||
// consume the single burst token
|
||||
// consume the single burst token
|
||||
let _ = rl.enqueue("t", b"x".to_vec()).await;
|
||||
// next should be rate limited (no refill)
|
||||
let result = rl.enqueue("t", b"y".to_vec()).await;
|
||||
|
||||
@@ -3,8 +3,8 @@
|
||||
use async_trait::async_trait;
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::error::{JobError, QueueError};
|
||||
use crate::job::Job;
|
||||
use crate::error::{QueueError, JobError};
|
||||
|
||||
/// A handle to a single unit of queued work.
|
||||
///
|
||||
|
||||
@@ -71,10 +71,13 @@ pub struct WorkerPool<Q: Queue + 'static> {
|
||||
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()
|
||||
})
|
||||
Self::with_config(
|
||||
queue,
|
||||
WorkerConfig {
|
||||
concurrency,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
/// Creates a new worker pool with explicit configuration.
|
||||
@@ -148,8 +151,8 @@ impl<Q: Queue + 'static> WorkerPool<Q> {
|
||||
/// 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 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)
|
||||
}
|
||||
|
||||
@@ -6,8 +6,8 @@
|
||||
//! service faster than its rate limit allows.
|
||||
|
||||
use crate::rate_limited::RateLimitedQueue;
|
||||
use crate::worker::{JobHandler, WorkerConfig, WorkerPool};
|
||||
use crate::traits::Queue;
|
||||
use crate::worker::{JobHandler, WorkerConfig, WorkerPool};
|
||||
|
||||
/// A `WorkerPool` whose dequeue is rate-limited via a token bucket.
|
||||
pub struct RateLimitedWorkerPool<Q: Queue + 'static> {
|
||||
@@ -40,12 +40,8 @@ mod tests {
|
||||
#[test]
|
||||
fn constructs_rate_limited_pool() {
|
||||
use crate::in_memory::InMemoryQueue;
|
||||
let _pool = RateLimitedWorkerPool::new(
|
||||
InMemoryQueue::new(),
|
||||
WorkerConfig::default(),
|
||||
10.0,
|
||||
5,
|
||||
);
|
||||
let _pool =
|
||||
RateLimitedWorkerPool::new(InMemoryQueue::new(), WorkerConfig::default(), 10.0, 5);
|
||||
// smoke: just verifies construction
|
||||
}
|
||||
}
|
||||
|
||||
@@ -36,7 +36,11 @@ async fn rate_limited_queue_no_item_loss_under_contention() {
|
||||
let s = Arc::clone(&seen);
|
||||
workers.push(tokio::spawn(async move {
|
||||
loop {
|
||||
match q.dequeue("stress", Duration::from_millis(50)).await.unwrap() {
|
||||
match q
|
||||
.dequeue("stress", Duration::from_millis(50))
|
||||
.await
|
||||
.unwrap()
|
||||
{
|
||||
Some(job) => {
|
||||
let _ = String::from_utf8(job.payload).unwrap();
|
||||
s.fetch_add(1, Ordering::SeqCst);
|
||||
@@ -64,4 +68,4 @@ async fn rate_limited_queue_no_item_loss_under_contention() {
|
||||
TASKS * PER_TASK,
|
||||
"items lost under contention"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user