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.
48 lines
1.6 KiB
Rust
48 lines
1.6 KiB
Rust
//! Rate-limited worker pool (feature `in-memory`).
|
|
//!
|
|
//! [`RateLimitedWorkerPool`] wraps [`crate::worker::WorkerPool`] with a
|
|
//! [`crate::rate_limited::RateLimitedQueue`] to back-pressure dequeue when the
|
|
//! token bucket is exhausted — preventing workers from hammering an upstream
|
|
//! service faster than its rate limit allows.
|
|
|
|
use crate::rate_limited::RateLimitedQueue;
|
|
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> {
|
|
inner: WorkerPool<RateLimitedQueue<Q>>,
|
|
}
|
|
|
|
impl<Q: Queue + 'static> RateLimitedWorkerPool<Q> {
|
|
/// Creates a rate-limited worker pool wrapping `queue` with the given
|
|
/// token-bucket rate (tokens/sec) and burst capacity.
|
|
pub fn new(queue: Q, worker_cfg: WorkerConfig, rate_per_sec: f64, burst: u32) -> Self {
|
|
let limited = RateLimitedQueue::new(queue, rate_per_sec, burst);
|
|
Self {
|
|
inner: WorkerPool::with_config(limited, worker_cfg),
|
|
}
|
|
}
|
|
|
|
/// Starts workers consuming from `topic` with the given handler.
|
|
pub fn start<H>(&self, topic: &str, handler: H)
|
|
where
|
|
H: JobHandler + 'static,
|
|
{
|
|
self.inner.start(topic, handler);
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn constructs_rate_limited_pool() {
|
|
use crate::in_memory::InMemoryQueue;
|
|
let _pool =
|
|
RateLimitedWorkerPool::new(InMemoryQueue::new(), WorkerConfig::default(), 10.0, 5);
|
|
// smoke: just verifies construction
|
|
}
|
|
}
|