pub mod config;
pub mod dlq;
pub mod error;
pub mod event;
pub mod job;
pub mod manager;
pub mod queue;
pub mod retry;
pub mod storage;
#[cfg(feature = "monitoring")]
pub mod alerts;
#[cfg(feature = "distributed")]
pub mod boost;
#[cfg(feature = "distributed")]
pub mod distributed;
#[cfg(feature = "metrics")]
pub mod metrics;
#[cfg(feature = "monitoring")]
pub mod monitor;
#[cfg(feature = "distributed")]
pub mod partition;
#[cfg(feature = "distributed")]
pub mod ratelimit;
#[cfg(feature = "telemetry")]
pub mod telemetry;
pub use config::LaneConfig;
pub use dlq::{DeadLetter, DeadLetterQueue};
pub use error::{LaneError, Result};
pub use event::{EventEmitter, EventPayload, EventStream, LaneEvent};
#[cfg(feature = "redis-backend")]
pub use job::RedisJobQueue;
pub use job::{
job_processor_fn, DeduplicationOptions, InMemoryJobQueue, Job, JobContext, JobEvent,
JobFinishedResult, JobFlow, JobFlowDependencyCountOptions, JobFlowDependencyCounts,
JobFlowDependencyKind, JobFlowDependencyPage, JobFlowDependencyPageCursor,
JobFlowDependencyPageItem, JobFlowDependencyPageOptions, JobFlowDependencyPages,
JobFlowDependencyPagesOptions, JobFlowDependencySelectedCounts, JobFlowDependencyValues, JobId,
JobLeaseRenewal, JobListOptions, JobListPage, JobLockToken, JobLogEntry, JobMetrics,
JobMetricsMeta, JobOptions, JobPriority, JobPriorityCount, JobProcessor, JobProcessorFn,
JobProcessorRouter, JobQueueBackend, JobQueueSnapshot, JobQueueStats, JobRateLimit,
JobRepeatEntry, JobRepeatListOptions, JobRepeatPage, JobRetention, JobRunOutcome, JobSpec,
JobState, JobStateCount, JobWorker, JobWorkerConfig, JobWorkerHandle, JobWorkerId,
LocalJobQueue, QueueName, RepeatOptions, RepeatSchedule, DEFAULT_JOB_EVENT_RETENTION,
DEFAULT_JOB_METRICS_RETENTION, DEFAULT_JOB_PRIORITY, MAX_JOB_PRIORITY,
};
pub use manager::{QueueManager, QueueManagerBuilder};
pub use queue::{
lane_ids, priorities, Command, CommandId, CommandQueue, JsonCommand, Lane, LaneId, LaneStatus,
Priority,
};
pub use retry::RetryPolicy;
pub use storage::{LocalStorage, Storage, StoredCommand, StoredDeadLetter};
#[cfg(feature = "monitoring")]
pub use alerts::{Alert, AlertLevel, AlertManager, LatencyAlertConfig, QueueDepthAlertConfig};
#[cfg(feature = "distributed")]
pub use boost::{PriorityBoostConfig, PriorityBooster};
#[cfg(feature = "distributed")]
pub use distributed::{
CommandEnvelope, CommandResult, DistributedQueue, LocalDistributedQueue, WorkerId, WorkerPool,
};
#[cfg(feature = "metrics")]
pub use metrics::{
metric_names, HistogramPercentiles, HistogramStats, LocalMetrics, MetricsBackend,
MetricsSnapshot, QueueMetrics,
};
#[cfg(feature = "monitoring")]
pub use monitor::{MonitorConfig, QueueMonitor};
#[cfg(feature = "distributed")]
pub use partition::{
CustomPartitioner, HashPartitioner, PartitionConfig, PartitionId, PartitionStrategy,
Partitioner, RoundRobinPartitioner,
};
#[cfg(feature = "distributed")]
pub use ratelimit::{RateLimitConfig, RateLimiter, SlidingWindowLimiter, TokenBucketLimiter};
#[cfg(feature = "telemetry")]
pub use telemetry::OtelMetricsBackend;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct QueueStats {
pub total_pending: usize,
pub total_active: usize,
pub dead_letter_count: usize,
pub lanes: HashMap<String, LaneStatus>,
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn test_queue_manager_builder() {
let emitter = EventEmitter::new(100);
let manager = QueueManagerBuilder::new(emitter)
.with_default_lanes()
.build()
.await
.unwrap();
let stats = manager.stats().await.unwrap();
assert_eq!(stats.lanes.len(), 6);
}
#[test]
fn test_queue_stats_default() {
let stats = QueueStats::default();
assert_eq!(stats.total_pending, 0);
assert_eq!(stats.total_active, 0);
assert!(stats.lanes.is_empty());
}
#[test]
fn test_queue_stats_serialization() {
let mut lanes = HashMap::new();
lanes.insert(
"query".to_string(),
LaneStatus {
pending: 5,
active: 2,
min: 1,
max: 10,
},
);
let stats = QueueStats {
total_pending: 5,
total_active: 2,
dead_letter_count: 0,
lanes,
};
let json = serde_json::to_string(&stats).unwrap();
let parsed: QueueStats = serde_json::from_str(&json).unwrap();
assert_eq!(parsed.total_pending, 5);
assert_eq!(parsed.total_active, 2);
}
}