# Fairness Queues Source: https://densumesh-broccoli-27.localhost:3000/advanced/fairness Ensure fair processing across message categories ## Overview Fairness queues prevent one category of messages from monopolizing processing resources. By using disambiguators, you can ensure messages are processed fairly across different tenants, users, or categories. Fairness queues are currently only supported with the **Redis** broker. ## The problem Without fairness, a burst of messages from one source can delay others: ```mermaid theme={null} timeline title Without Fairness section Queue Tenant A : 100 messages Tenant B : 5 messages Tenant C : 5 messages section Processing 0-100ms : All Tenant A 100-105ms : Tenant B 105-110ms : Tenant C ``` Tenant B and C must wait for all of Tenant A's messages to process. ## With fairness Fairness distributes processing across categories: ```mermaid theme={null} timeline title With Fairness section Processing 0-3ms : A, B, C (round robin) 3-6ms : A, B, C ... : continues fairly ``` ## Using disambiguators ### Publishing with disambiguators Include a disambiguator when publishing: ```rust theme={null} // Each tenant gets their own sub-queue queue.publish("jobs", Some("tenant-a".into()), &job_a, None).await?; queue.publish("jobs", Some("tenant-b".into()), &job_b, None).await?; queue.publish("jobs", Some("tenant-c".into()), &job_c, None).await?; ``` ### Consuming with fairness Enable fairness in consume options: ```rust theme={null} use broccoli_queue::queue::ConsumeOptions; let options = ConsumeOptions::builder() .fairness(true) .build(); queue.process_messages("jobs", Some(4), Some(options), |msg| async move { println!("Processing for {:?}", msg.disambiguator); Ok(()) }).await?; ``` ## How it works With fairness enabled: 1. Messages are stored in sub-queues based on disambiguator 2. The consumer rotates through sub-queues in round-robin fashion 3. Each consume operation pulls from the next sub-queue ```mermaid theme={null} flowchart LR subgraph "Main Queue: jobs" direction TB A[tenant-a queue] B[tenant-b queue] C[tenant-c queue] end Consumer --> |round robin| A Consumer --> |round robin| B Consumer --> |round robin| C ``` ## Use cases ### Multi-tenant applications Prevent one tenant from affecting others: ```rust theme={null} async fn handle_tenant_job(tenant_id: &str, job: TenantJob) { queue.publish( "tenant-jobs", Some(tenant_id.to_string()), &job, None ).await?; } // Consumer processes fairly across tenants let options = ConsumeOptions::builder().fairness(true).build(); queue.process_messages("tenant-jobs", Some(8), Some(options), handler).await?; ``` ### Priority levels Implement soft priority with multiple queues: ```rust theme={null} // High priority gets its own queue queue.publish("jobs", Some("priority-high".into()), &urgent_job, None).await?; // Normal priority queue.publish("jobs", Some("priority-normal".into()), &normal_job, None).await?; // With fairness, both get equal attention // For true priority, use separate queues or priority options ``` ### User-based fairness Prevent a single user from overwhelming the system: ```rust theme={null} let user_id = request.user_id; queue.publish("user-jobs", Some(user_id), &job, None).await?; ``` ## Queue size with fairness When using fairness, `queue.size()` returns sizes for each sub-queue: ```rust theme={null} let sizes = queue.size("jobs").await?; for (queue_name, size) in sizes { println!("{}: {} messages", queue_name, size); } // Output: // jobs:fairness:tenant-a: 45 messages // jobs:fairness:tenant-b: 12 messages // jobs:fairness:tenant-c: 8 messages ``` ## Configuration options ```rust theme={null} let options = ConsumeOptions::builder() .fairness(true) .auto_ack(false) // Manual acknowledgment .consume_wait(Duration::from_millis(10)) // Wait between iterations .build(); ``` ## Best practices Use identifiers that represent your fairness boundaries: * Tenant ID for multi-tenant apps * User ID for user fairness * Region for geographic distribution * Priority level for soft prioritization Track the size of each sub-queue to detect imbalances: ```rust theme={null} let sizes = queue.size("jobs").await?; for (name, size) in sizes { metrics::gauge!("queue.size", size as f64, "queue" => name); } ``` Too many unique disambiguators can impact performance. Consider bucketing: ```rust theme={null} // Instead of per-user, bucket by user hash let bucket = format!("bucket-{}", user_id.hash() % 100); queue.publish("jobs", Some(bucket), &job, None).await?; ``` ## Limitations * Only available with Redis broker * Round-robin is fixed (no weighted fairness) * Empty sub-queues are checked (slight overhead) ## Alternatives If fairness queues don't fit your needs: 1. **Separate queues**: Create distinct queues per category 2. **Priority option**: Use `PublishOptions::priority()` for priority-based ordering 3. **Custom routing**: Implement your own routing logic with multiple queues # Management API Source: https://densumesh-broccoli-27.localhost:3000/advanced/management Monitor and manage your queues ## Overview The management feature provides APIs to inspect queue status and monitor your message processing system. ## Enabling management Add the `management` feature to your `Cargo.toml`: ```toml theme={null} [dependencies] broccoli_queue = { version = "0.4", features = ["management"] } ``` Management is currently supported for Redis and RabbitMQ brokers only. ## Queue status Get the status of a queue: ```rust theme={null} #[cfg(feature = "management")] async fn check_queue_health(queue: &BroccoliQueue) -> Result<(), BroccoliError> { let status = queue.queue_status("jobs".into(), None).await?; println!("Queue: jobs"); println!(" Pending: {}", status.pending_count); println!(" Processing: {}", status.processing_count); println!(" Failed: {}", status.failed_count); Ok(()) } ``` ### With disambiguator For fairness queues, specify the disambiguator: ```rust theme={null} let status = queue.queue_status( "jobs".into(), Some("tenant-123".into()) ).await?; ``` ## Queue size Get queue sizes (available without management feature): ```rust theme={null} let sizes = queue.size("jobs").await?; for (queue_name, size) in sizes { println!("{}: {} messages", queue_name, size); } ``` For fairness queues, this returns sizes for each sub-queue. ## Building a monitoring dashboard Example of exposing queue metrics: ```rust theme={null} use std::collections::HashMap; #[derive(Serialize)] struct QueueMetrics { name: String, pending: u64, processing: u64, failed: u64, } async fn get_metrics(queue: &BroccoliQueue, queue_names: &[&str]) -> Vec { let mut metrics = Vec::new(); for name in queue_names { #[cfg(feature = "management")] if let Ok(status) = queue.queue_status(name.to_string(), None).await { metrics.push(QueueMetrics { name: name.to_string(), pending: status.pending_count, processing: status.processing_count, failed: status.failed_count, }); } } metrics } ``` ## Health checks Implement health checks for your queue system: ```rust theme={null} async fn health_check(queue: &BroccoliQueue) -> bool { // Try to get queue size as a connectivity check match queue.size("health-check").await { Ok(_) => true, Err(e) => { log::error!("Queue health check failed: {:?}", e); false } } } ``` ## Alerting on failed messages Monitor the failed queue for alerting: ```rust theme={null} #[cfg(feature = "management")] async fn check_failed_messages(queue: &BroccoliQueue) -> Result<(), BroccoliError> { let status = queue.queue_status("jobs".into(), None).await?; if status.failed_count > 100 { // Send alert alert::send("High failed message count", &format!( "Queue 'jobs' has {} failed messages", status.failed_count )).await; } Ok(()) } ``` ## Metrics integration ### Prometheus ```rust theme={null} use prometheus::{IntGauge, register_int_gauge}; lazy_static! { static ref QUEUE_PENDING: IntGauge = register_int_gauge!( "broccoli_queue_pending", "Number of pending messages" ).unwrap(); static ref QUEUE_PROCESSING: IntGauge = register_int_gauge!( "broccoli_queue_processing", "Number of messages being processed" ).unwrap(); static ref QUEUE_FAILED: IntGauge = register_int_gauge!( "broccoli_queue_failed", "Number of failed messages" ).unwrap(); } #[cfg(feature = "management")] async fn update_metrics(queue: &BroccoliQueue) { if let Ok(status) = queue.queue_status("jobs".into(), None).await { QUEUE_PENDING.set(status.pending_count as i64); QUEUE_PROCESSING.set(status.processing_count as i64); QUEUE_FAILED.set(status.failed_count as i64); } } ``` ### OpenTelemetry ```rust theme={null} use opentelemetry::metrics::Meter; async fn record_metrics(meter: &Meter, queue: &BroccoliQueue) { let pending_gauge = meter.i64_gauge("broccoli.queue.pending").init(); let failed_gauge = meter.i64_gauge("broccoli.queue.failed").init(); #[cfg(feature = "management")] if let Ok(status) = queue.queue_status("jobs".into(), None).await { pending_gauge.record(status.pending_count as i64, &[]); failed_gauge.record(status.failed_count as i64, &[]); } } ``` ## Redis CLI inspection For Redis, you can also inspect queues directly: ```bash theme={null} # Main queue size redis-cli LLEN jobs # Processing queue redis-cli ZCARD jobs:processing # Failed queue redis-cli LLEN jobs:failed # Scheduled messages redis-cli ZCARD jobs:scheduled # List failed messages redis-cli LRANGE jobs:failed 0 10 ``` ## Best practices Don't query status on every request. Use a background task: ```rust theme={null} tokio::spawn(async move { let mut interval = tokio::time::interval(Duration::from_secs(30)); loop { interval.tick().await; update_metrics(&queue).await; } }); ``` Monitor the failed queue and alert when it exceeds thresholds. Measure how long messages spend in processing state to detect stuck jobs. # Message Scheduling Source: https://densumesh-broccoli-27.localhost:3000/advanced/scheduling Schedule messages for delayed or future delivery ## Overview Broccoli supports scheduling messages for delayed delivery. You can either delay a message by a duration or schedule it for a specific time. ## Enabling scheduling Enable scheduling when building your queue: ```rust theme={null} let queue = BroccoliQueue::builder("redis://localhost:6379") .enable_scheduling(true) .build() .await?; ``` For RabbitMQ, you must install the [delayed-exchange plugin](https://www.rabbitmq.com/blog/2015/04/16/scheduling-messages-with-rabbitmq). ## Delayed delivery Delay a message by a specific duration: ```rust theme={null} use broccoli_queue::queue::PublishOptions; use time::Duration; let options = PublishOptions::builder() .delay(Duration::seconds(30)) .build(); queue.publish("jobs", None, &job, Some(options)).await?; ``` The message will be delivered approximately 30 seconds after publishing. ## Scheduled delivery Schedule a message for a specific time: ```rust theme={null} use broccoli_queue::queue::PublishOptions; use time::OffsetDateTime; // Schedule for tomorrow at midnight let scheduled_time = OffsetDateTime::now_utc() .replace_time(time::Time::MIDNIGHT) + time::Duration::days(1); let options = PublishOptions::builder() .schedule_at(scheduled_time) .build(); queue.publish("jobs", None, &job, Some(options)).await?; ``` ## Batch scheduling Schedule multiple messages: ```rust theme={null} use time::Duration; let jobs = vec![ JobPayload { id: "1".into(), task: "task-a".into() }, JobPayload { id: "2".into(), task: "task-b".into() }, ]; let options = PublishOptions::builder() .delay(Duration::minutes(5)) .build(); queue.publish_batch("jobs", None, jobs, Some(options)).await?; ``` ## PublishOptions The `PublishOptions` struct supports several scheduling-related options: ```rust theme={null} pub struct PublishOptions { /// Time-to-live for the message pub ttl: Option, /// Message priority (1-5, where 1 is highest) pub priority: Option, /// Delay before delivery pub delay: Option, /// Specific delivery time pub scheduled_at: Option, } ``` ### Builder pattern ```rust theme={null} let options = PublishOptions::builder() .delay(Duration::seconds(60)) .priority(1) // High priority .ttl(Duration::hours(24)) .build(); ``` ## Use cases ### Scheduled reports ```rust theme={null} #[derive(Clone, Serialize, Deserialize)] struct ReportJob { report_type: String, recipients: Vec, } // Schedule daily report for 6 AM let tomorrow_6am = OffsetDateTime::now_utc() .replace_time(time::Time::from_hms(6, 0, 0)?) + time::Duration::days(1); let report = ReportJob { report_type: "daily_summary".into(), recipients: vec!["admin@example.com".into()], }; queue.publish( "reports", None, &report, Some(PublishOptions::builder().schedule_at(tomorrow_6am).build()) ).await?; ``` ### Retry with backoff ```rust theme={null} async fn handle_with_retry( queue: &BroccoliQueue, job: JobPayload, attempt: u8, ) -> Result<(), BroccoliError> { match process_job(&job).await { Ok(_) => Ok(()), Err(e) if attempt < 5 => { // Exponential backoff: 1s, 2s, 4s, 8s, 16s let delay = Duration::seconds(2_i64.pow(attempt as u32)); queue.publish( "jobs", None, &job, Some(PublishOptions::builder().delay(delay).build()) ).await?; Ok(()) } Err(e) => Err(e), } } ``` ### Rate limiting ```rust theme={null} // Spread 100 emails over 10 minutes to avoid rate limits for (i, email) in emails.iter().enumerate() { let delay = Duration::seconds((i as i64) * 6); // One every 6 seconds queue.publish( "emails", None, email, Some(PublishOptions::builder().delay(delay).build()) ).await?; } ``` ## How scheduling works ```mermaid theme={null} sequenceDiagram participant P as Producer participant B as Broker participant S as Scheduled Store participant Q as Main Queue participant C as Consumer P->>B: publish(delay=30s) B->>S: Store with scheduled_at Note over S: Wait until scheduled_at S->>Q: Move to main queue Q->>C: Consumer picks up message ``` ## Broker-specific behavior | Feature | Redis | RabbitMQ | SurrealDB | | -------------- | ---------- | ------------------- | ---------- | | `delay` | ✅ | ✅ (requires plugin) | ✅ | | `scheduled_at` | ✅ | ✅ (requires plugin) | ✅ | | Precision | \~1 second | \~1 second | \~1 second | ## Best practices For delays under a few hours, use `delay`: ```rust theme={null} PublishOptions::builder().delay(Duration::minutes(30)).build() ``` For calendar-based scheduling, use `scheduled_at`: ```rust theme={null} PublishOptions::builder().schedule_at(next_monday_9am).build() ``` `scheduled_at` uses `OffsetDateTime`. Be explicit about time zones: ```rust theme={null} use time::{OffsetDateTime, UtcOffset}; // Schedule for 9 AM in UTC-5 (EST) let est = UtcOffset::from_hms(-5, 0, 0)?; let scheduled = OffsetDateTime::now_utc() .to_offset(est) .replace_time(time::Time::from_hms(9, 0, 0)?); ``` With Redis, check the scheduled queue: ```bash theme={null} redis-cli ZCARD jobs:scheduled ``` # Errors Source: https://densumesh-broccoli-27.localhost:3000/api/errors API reference for BroccoliError and error handling ## BroccoliError The error type for all Broccoli operations. ```rust theme={null} use broccoli_queue::error::BroccoliError; ``` ### Variants ```rust theme={null} #[derive(Debug, thiserror::Error)] pub enum BroccoliError { /// Broker connection or operation errors #[error("Broker error: {0}")] Broker(String), /// Non-idempotent concurrent operation (SurrealDB) #[error("Broker error: non-idempotent operation {0}")] BrokerNonIdempotentOp(String), /// Retriable non-idempotent operation (SurrealDB) #[error("Broker error: non-idempotent retriable operation {0}")] BrokerNonIdempotentRetriableOp(String), /// Failed to publish message #[error("Failed to publish message: {0}")] Publish(String), /// Failed to consume message #[error("Failed to consume message: {0}")] Consume(String), /// Failed to acknowledge message #[error("Failed to acknowledge message: {0}")] Acknowledge(String), /// Failed to reject message #[error("Failed to reject message: {0}")] Reject(String), /// Failed to cancel message #[error("Failed to cancel message: {0}")] Cancel(String), /// Failed to get message position #[error("Failed to get message position: {0}")] GetMessagePosition(String), /// JSON serialization/deserialization error #[error("Serialization error: {0}")] Serialization(#[from] serde_json::Error), /// Redis-specific error (with redis feature) #[cfg(feature = "redis")] #[error("Redis error: {0}")] Redis(#[from] redis::RedisError), /// SurrealDB-specific error (with surrealdb feature) #[cfg(feature = "surrealdb")] #[error("SurrealDB error: {0}")] SurrealDB(#[from] surrealdb::Error), /// Job processing error #[error("Job error: {0}")] Job(String), /// Queue status retrieval error #[error("Queue status error: {0}")] QueueStatus(String), /// Connection timeout #[error("Connection timeout after {0} retries")] ConnectionTimeout(u32), /// Feature not implemented for broker #[error("Feature not implemented")] NotImplemented, } ``` ## Error handling ### Basic error handling ```rust theme={null} use broccoli_queue::error::BroccoliError; async fn publish_job(queue: &BroccoliQueue, job: &Job) -> Result<(), BroccoliError> { queue.publish("jobs", None, job, None).await?; Ok(()) } // Usage match publish_job(&queue, &job).await { Ok(_) => println!("Published successfully"), Err(BroccoliError::Publish(msg)) => eprintln!("Publish failed: {}", msg), Err(BroccoliError::Broker(msg)) => eprintln!("Broker error: {}", msg), Err(e) => eprintln!("Other error: {}", e), } ``` ### In message handlers ```rust theme={null} queue.process_messages("jobs", Some(4), None, |msg: BrokerMessage| async move { // Return BroccoliError::Job for business logic failures if !is_valid(&msg.payload) { return Err(BroccoliError::Job("Invalid job data".into())); } // Process the job process(&msg.payload).await.map_err(|e| { BroccoliError::Job(format!("Processing failed: {}", e)) })?; Ok(()) }).await?; ``` ### Converting from other errors ```rust theme={null} async fn process_job(job: &Job) -> Result<(), BroccoliError> { // From serde_json::Error let data: Value = serde_json::from_str(&job.data)?; // From custom errors external_service(&data) .await .map_err(|e| BroccoliError::Job(e.to_string()))?; Ok(()) } ``` ## Common error scenarios ### Connection errors ```rust theme={null} let result = BroccoliQueue::builder("redis://invalid:6379") .build() .await; match result { Err(BroccoliError::Broker(msg)) => { eprintln!("Failed to connect: {}", msg); // Handle reconnection or fallback } _ => {} } ``` ### Serialization errors ```rust theme={null} // This would fail if Job doesn't implement Serialize let result = queue.publish("jobs", None, &invalid_job, None).await; match result { Err(BroccoliError::Serialization(e)) => { eprintln!("Failed to serialize message: {}", e); } _ => {} } ``` ### Timeout errors ```rust theme={null} match result { Err(BroccoliError::ConnectionTimeout(retries)) => { eprintln!("Connection timed out after {} retries", retries); } _ => {} } ``` ### Feature not implemented ```rust theme={null} // Some operations aren't available on all brokers match queue.some_operation().await { Err(BroccoliError::NotImplemented) => { eprintln!("This feature isn't supported by your broker"); } _ => {} } ``` ## Error recovery patterns ### Retry with backoff ```rust theme={null} async fn publish_with_retry( queue: &BroccoliQueue, job: &Job, max_retries: u32, ) -> Result<(), BroccoliError> { let mut attempts = 0; loop { match queue.publish("jobs", None, job, None).await { Ok(_) => return Ok(()), Err(BroccoliError::Broker(_)) if attempts < max_retries => { attempts += 1; let delay = std::time::Duration::from_millis(100 * 2_u64.pow(attempts)); tokio::time::sleep(delay).await; } Err(e) => return Err(e), } } } ``` ### Graceful degradation ```rust theme={null} async fn process_with_fallback(job: &Job) -> Result<(), BroccoliError> { match primary_processing(job).await { Ok(_) => Ok(()), Err(BroccoliError::Job(_)) => { // Try fallback processing fallback_processing(job).await } Err(e) => Err(e), // Propagate other errors } } ``` ## Logging errors ```rust theme={null} use log::{error, warn}; queue.process_messages("jobs", Some(4), None, |msg: BrokerMessage| async move { match process_job(&msg.payload).await { Ok(_) => Ok(()), Err(e) => { error!( "Job {} failed (attempt {}): {:?}", msg.task_id, msg.attempts, e ); Err(e) } } }).await?; ``` ## Custom error types If you need richer error information in your handlers: ```rust theme={null} #[derive(Debug)] enum JobError { ValidationFailed(String), ExternalServiceError(String), ResourceNotFound(String), } impl From for BroccoliError { fn from(e: JobError) -> Self { BroccoliError::Job(format!("{:?}", e)) } } async fn process_job(job: &Job) -> Result<(), JobError> { if !is_valid(job) { return Err(JobError::ValidationFailed("Invalid fields".into())); } Ok(()) } // In handler queue.process_messages("jobs", Some(4), None, |msg| async move { process_job(&msg.payload).await.map_err(Into::into) }).await?; ``` # Messages Source: https://densumesh-broccoli-27.localhost:3000/api/messages API reference for BrokerMessage and related types ## BrokerMessage A wrapper for messages that includes metadata for processing. ```rust theme={null} use broccoli_queue::brokers::broker::BrokerMessage; ``` ### Structure ```rust theme={null} pub struct BrokerMessage { /// Unique identifier for the message pub task_id: uuid::Uuid, /// The actual message content pub payload: T, /// Number of processing attempts made pub attempts: u8, /// Disambiguator for message fairness pub disambiguator: Option, } ``` ### Fields | Field | Type | Description | | --------------- | ---------------- | ---------------------------------------------------- | | `task_id` | `Uuid` | Unique message identifier, auto-generated on publish | | `payload` | `T` | Your custom message data | | `attempts` | `u8` | Number of times this message has been processed | | `disambiguator` | `Option` | Optional identifier for fairness queue routing | ### Creating messages Messages are automatically created when publishing: ```rust theme={null} let job = JobPayload { id: "1".into(), task: "process".into() }; // publish() creates and returns the BrokerMessage let message = queue.publish("jobs", None, &job, None).await?; println!("Task ID: {}", message.task_id); ``` ### Accessing message data ```rust theme={null} queue.process_messages("jobs", Some(1), None, |msg: BrokerMessage| async move { // Access metadata println!("ID: {}", msg.task_id); println!("Attempts: {}", msg.attempts); println!("Disambiguator: {:?}", msg.disambiguator); // Access payload let job = &msg.payload; println!("Job: {} - {}", job.id, job.task); Ok(()) }).await?; ``` *** ## Payload requirements Your message payload type must implement: * `Clone` * `serde::Serialize` * `serde::Deserialize` ### Example payload ```rust theme={null} use serde::{Serialize, Deserialize}; #[derive(Debug, Clone, Serialize, Deserialize)] pub struct JobPayload { pub id: String, pub task_name: String, pub parameters: serde_json::Value, pub created_at: chrono::DateTime, } ``` ### Complex payload ```rust theme={null} #[derive(Debug, Clone, Serialize, Deserialize)] pub struct EmailJob { pub recipient: String, pub subject: String, pub body: String, pub attachments: Vec, } #[derive(Debug, Clone, Serialize, Deserialize)] pub struct Attachment { pub filename: String, pub content_type: String, #[serde(with = "base64_serde")] pub data: Vec, } ``` *** ## InternalBrokerMessage Internal message representation used by broker implementations. You typically don't interact with this directly. ```rust theme={null} pub struct InternalBrokerMessage { pub task_id: String, pub payload: String, // JSON serialized pub attempts: u8, pub disambiguator: Option, } ``` *** ## BrokerConfig Configuration options for broker behavior. ```rust theme={null} pub struct BrokerConfig { /// Maximum retry attempts (default: 3) pub retry_attempts: Option, /// Whether to retry failed messages (default: true) pub retry_failed: Option, /// Connection pool size (default: 10) pub pool_connections: Option, /// Enable message scheduling (default: false) pub enable_scheduling: Option, } ``` ### Default values ```rust theme={null} impl Default for BrokerConfig { fn default() -> Self { Self { retry_attempts: Some(3), retry_failed: Some(true), pool_connections: Some(10), enable_scheduling: Some(false), } } } ``` *** ## BrokerType Enum representing supported broker types. ```rust theme={null} pub enum BrokerType { #[cfg(feature = "redis")] Redis, #[cfg(feature = "rabbitmq")] RabbitMQ, #[cfg(feature = "surrealdb")] SurrealDB, } ``` *** ## Broker trait The `Broker` trait defines the interface that all broker implementations must satisfy. This is internal to Broccoli but useful for understanding the abstraction. ```rust theme={null} #[async_trait] pub trait Broker: Send + Sync { async fn connect(&mut self, broker_url: &str) -> Result<(), BroccoliError>; async fn publish( &self, queue_name: &str, disambiguator: Option, message: &[InternalBrokerMessage], options: Option, ) -> Result, BroccoliError>; async fn consume( &self, queue_name: &str, options: Option, ) -> Result; async fn try_consume( &self, queue_name: &str, options: Option, ) -> Result, BroccoliError>; async fn acknowledge( &self, queue_name: &str, message: InternalBrokerMessage, ) -> Result<(), BroccoliError>; async fn reject( &self, queue_name: &str, message: InternalBrokerMessage, ) -> Result<(), BroccoliError>; async fn cancel( &self, queue_name: &str, message_id: String, ) -> Result<(), BroccoliError>; async fn size( &self, queue_name: &str, ) -> Result, BroccoliError>; } ``` # Options Source: https://densumesh-broccoli-27.localhost:3000/api/options API reference for PublishOptions, ConsumeOptions, and RetryStrategy ## PublishOptions Options for publishing messages. ```rust theme={null} use broccoli_queue::queue::PublishOptions; ``` ### Structure ```rust theme={null} pub struct PublishOptions { /// Time-to-live for the message pub ttl: Option, /// Message priority (1-5, where 1 is highest) pub priority: Option, /// Delay before the message is published pub delay: Option, /// Scheduled time for message delivery pub scheduled_at: Option, } ``` ### Builder ```rust theme={null} let options = PublishOptions::builder() .delay(Duration::seconds(30)) .priority(1) .ttl(Duration::hours(24)) .build(); queue.publish("jobs", None, &job, Some(options)).await?; ``` ### Methods #### `builder()` Creates a new `PublishOptionsBuilder`. ```rust theme={null} pub const fn builder() -> PublishOptionsBuilder ``` *** ## PublishOptionsBuilder Builder for constructing `PublishOptions`. ### `ttl` Sets the time-to-live for the message. ```rust theme={null} pub const fn ttl(mut self, duration: Duration) -> Self ``` **Example:** ```rust theme={null} PublishOptions::builder() .ttl(Duration::hours(1)) .build() ``` ### `priority` Sets the priority level (1-5, where 1 is highest). ```rust theme={null} pub fn priority(mut self, priority: u8) -> Self ``` **Panics:** If priority is not between 1 and 5. **Example:** ```rust theme={null} PublishOptions::builder() .priority(1) // High priority .build() ``` ### `delay` Sets a delay before message delivery. ```rust theme={null} pub const fn delay(mut self, duration: Duration) -> Self ``` **Example:** ```rust theme={null} PublishOptions::builder() .delay(Duration::minutes(5)) .build() ``` ### `schedule_at` Sets a specific delivery time. ```rust theme={null} pub const fn schedule_at(mut self, time: OffsetDateTime) -> Self ``` **Example:** ```rust theme={null} use time::OffsetDateTime; let tomorrow = OffsetDateTime::now_utc() + Duration::days(1); PublishOptions::builder() .schedule_at(tomorrow) .build() ``` ### `build` Builds the `PublishOptions`. ```rust theme={null} pub const fn build(self) -> PublishOptions ``` *** ## ConsumeOptions Options for consuming messages. ```rust theme={null} use broccoli_queue::queue::ConsumeOptions; ``` ### Structure ```rust theme={null} pub struct ConsumeOptions { /// Auto-acknowledge messages (default: false) pub auto_ack: Option, /// Enable fairness queue consumption (Redis only) pub fairness: Option, /// Wait duration between consume iterations pub consume_wait: Option, /// Acknowledge after handler success (default: true) pub handler_ack: Option, } ``` ### Builder ```rust theme={null} let options = ConsumeOptions::builder() .fairness(true) .auto_ack(false) .build(); queue.process_messages("jobs", Some(4), Some(options), handler).await?; ``` *** ## ConsumeOptionsBuilder Builder for constructing `ConsumeOptions`. ### `auto_ack` Sets whether messages are auto-acknowledged. ```rust theme={null} pub const fn auto_ack(mut self, auto_ack: bool) -> Self ``` **Default:** `false` If `auto_ack` is true, calling `acknowledge()` or `reject()` will return an error. ### `fairness` Enables fairness queue consumption (Redis only). ```rust theme={null} pub const fn fairness(mut self, fairness: bool) -> Self ``` **Example:** ```rust theme={null} ConsumeOptions::builder() .fairness(true) .build() ``` ### `consume_wait` Sets the wait duration between consume loop iterations. ```rust theme={null} pub const fn consume_wait(mut self, consume_wait: std::time::Duration) -> Self ``` This allows consumer loops to be interrupted by tokio. **Example:** ```rust theme={null} ConsumeOptions::builder() .consume_wait(std::time::Duration::from_millis(10)) .build() ``` ### `handler_ack` Controls automatic acknowledgment after successful handler execution. ```rust theme={null} pub const fn handler_ack(mut self, followup: bool) -> Self ``` **Default:** `true` Set to `false` to manually control acknowledgment: ```rust theme={null} ConsumeOptions::builder() .handler_ack(false) .build() ``` ### `build` Builds the `ConsumeOptions`. ```rust theme={null} pub const fn build(self) -> ConsumeOptions ``` *** ## RetryStrategy Configuration for message retry behavior. ```rust theme={null} use broccoli_queue::queue::RetryStrategy; ``` ### Structure ```rust theme={null} pub struct RetryStrategy { /// Whether failed messages should be retried pub retry_failed: bool, /// Maximum number of retry attempts pub attempts: Option, } ``` ### Default ```rust theme={null} impl Default for RetryStrategy { fn default() -> Self { Self { retry_failed: true, attempts: Some(3), } } } ``` ### Methods #### `new` Creates a new retry strategy with defaults. ```rust theme={null} pub const fn new() -> Self ``` #### `with_attempts` Sets the maximum retry attempts. ```rust theme={null} pub const fn with_attempts(mut self, attempts: u8) -> Self ``` **Example:** ```rust theme={null} RetryStrategy::new().with_attempts(5) ``` #### `retry_failed` Enables or disables retrying. ```rust theme={null} pub const fn retry_failed(mut self, retry_failed: bool) -> Self ``` **Example:** ```rust theme={null} // Disable retries - failures go directly to failed queue RetryStrategy::new().retry_failed(false) ``` ### Usage ```rust theme={null} let queue = BroccoliQueue::builder("redis://localhost:6379") .failed_message_retry_strategy( RetryStrategy::new() .with_attempts(5) .retry_failed(true) ) .build() .await?; ``` *** ## QueueStatus (management feature) Status information for a queue. ```rust theme={null} #[cfg(feature = "management")] use broccoli_queue::brokers::management::QueueStatus; ``` ### Structure ```rust theme={null} pub struct QueueStatus { pub pending_count: u64, pub processing_count: u64, pub failed_count: u64, } ``` ### Usage ```rust theme={null} #[cfg(feature = "management")] { let status = queue.queue_status("jobs".into(), None).await?; println!("Pending: {}", status.pending_count); println!("Processing: {}", status.processing_count); println!("Failed: {}", status.failed_count); } ``` # BroccoliQueue Source: https://densumesh-broccoli-27.localhost:3000/api/queue API reference for the main queue interface ## BroccoliQueue The main entry point for interacting with Broccoli message queues. ```rust theme={null} use broccoli_queue::queue::BroccoliQueue; ``` ## Creating a queue ### `builder` Creates a new `BroccoliQueueBuilder` for configuring the queue. ```rust theme={null} pub fn builder(broker_url: impl Into) -> BroccoliQueueBuilder ``` **Parameters:** * `broker_url` - Connection URL for the message broker **Example:** ```rust theme={null} let queue = BroccoliQueue::builder("redis://localhost:6379") .pool_connections(5) .build() .await?; ``` ### `builder_with` (SurrealDB only) Creates a builder with an existing SurrealDB connection. ```rust theme={null} #[cfg(feature = "surrealdb")] pub fn builder_with(db: Surreal) -> BroccoliQueueBuilder ``` ## Publishing ### `publish` Publishes a single message to a queue. ```rust theme={null} pub async fn publish( &self, topic: &str, disambiguator: Option, message: &T, options: Option, ) -> Result, BroccoliError> where T: Clone + Serialize + DeserializeOwned ``` **Parameters:** * `topic` - Queue name * `disambiguator` - Optional identifier for fairness queues * `message` - The message payload * `options` - Optional publish configuration **Returns:** The wrapped `BrokerMessage` with generated `task_id` **Example:** ```rust theme={null} let message = queue.publish("jobs", None, &job, None).await?; println!("Published: {}", message.task_id); ``` ### `publish_batch` Publishes multiple messages to a queue. ```rust theme={null} pub async fn publish_batch( &self, topic: &str, disambiguator: Option, messages: impl IntoIterator, options: Option, ) -> Result>, BroccoliError> where T: Clone + Serialize + DeserializeOwned ``` **Example:** ```rust theme={null} let jobs = vec![job1, job2, job3]; let messages = queue.publish_batch("jobs", None, jobs, None).await?; ``` ## Consuming ### `consume` Consumes a message, blocking until one is available. ```rust theme={null} pub async fn consume( &self, topic: &str, options: Option, ) -> Result, BroccoliError> where T: Clone + Serialize + DeserializeOwned ``` **Example:** ```rust theme={null} let message: BrokerMessage = queue.consume("jobs", None).await?; ``` ### `try_consume` Attempts to consume a message without blocking. ```rust theme={null} pub async fn try_consume( &self, topic: &str, options: Option, ) -> Result>, BroccoliError> where T: Clone + Serialize + DeserializeOwned ``` **Returns:** `Some(message)` if available, `None` otherwise **Example:** ```rust theme={null} if let Some(message) = queue.try_consume::("jobs", None).await? { // Process message } ``` ### `consume_batch` Consumes multiple messages with a timeout. ```rust theme={null} pub async fn consume_batch( &self, topic: &str, batch_size: usize, timeout: Duration, options: Option, ) -> Result>, BroccoliError> where T: Clone + Serialize + DeserializeOwned ``` **Example:** ```rust theme={null} let messages = queue.consume_batch::( "jobs", 10, Duration::seconds(5), None ).await?; ``` ### `try_consume_batch` Attempts to consume multiple messages without blocking. ```rust theme={null} pub async fn try_consume_batch( &self, topic: &str, batch_size: usize, options: Option, ) -> Result>, BroccoliError> ``` ## Message handling ### `acknowledge` Acknowledges successful processing of a message. ```rust theme={null} pub async fn acknowledge( &self, topic: &str, message: BrokerMessage, ) -> Result<(), BroccoliError> where T: Clone + Serialize + DeserializeOwned ``` **Example:** ```rust theme={null} queue.acknowledge("jobs", message).await?; ``` ### `reject` Rejects a message, triggering retry or failure handling. ```rust theme={null} pub async fn reject( &self, topic: &str, message: BrokerMessage, ) -> Result<(), BroccoliError> where T: Clone + Serialize + DeserializeOwned ``` ### `cancel` Cancels a message by ID. ```rust theme={null} pub async fn cancel( &self, topic: &str, message_id: String, ) -> Result<(), BroccoliError> ``` ## Processing ### `process_messages` Processes messages in a loop with a handler function. ```rust theme={null} pub async fn process_messages( &self, topic: &str, concurrency: Option, consume_options: Option, handler: F, ) -> Result<(), BroccoliError> where T: DeserializeOwned + Send + Clone + Serialize + 'static, F: Fn(BrokerMessage) -> Fut + Send + Sync + Clone + 'static, Fut: Future> + Send + 'static ``` **Parameters:** * `topic` - Queue name * `concurrency` - Number of concurrent workers (None for single-threaded) * `consume_options` - Optional consume configuration * `handler` - Async function to process each message **Example:** ```rust theme={null} queue.process_messages("jobs", Some(4), None, |msg: BrokerMessage| async move { println!("Processing: {:?}", msg.payload); Ok(()) }).await?; ``` ### `process_messages_with_handlers` Processes messages with separate success and error handlers. ```rust theme={null} pub async fn process_messages_with_handlers( &self, topic: &str, concurrency: Option, consume_options: Option, message_handler: F, on_success: S, on_error: E, ) -> Result<(), BroccoliError> ``` **Example:** ```rust theme={null} queue.process_messages_with_handlers( "jobs", Some(4), None, |msg| async move { process(msg.payload).await }, |msg, result| async move { log_success(&msg).await; Ok(()) }, |msg, error| async move { log_error(&msg, &error).await; Ok(()) }, ).await?; ``` ## Queue information ### `size` Returns the size of queue(s). ```rust theme={null} pub async fn size( &self, queue_name: &str, ) -> Result, BroccoliError> ``` **Returns:** Map of queue names to message counts ### `queue_status` (management feature) Returns detailed queue status. ```rust theme={null} #[cfg(feature = "management")] pub async fn queue_status( &self, queue_name: String, disambiguator: Option, ) -> Result ``` *** ## BroccoliQueueBuilder Builder for configuring `BroccoliQueue`. ### `pool_connections` Sets the connection pool size. ```rust theme={null} pub const fn pool_connections(mut self, connections: u8) -> Self ``` **Default:** 10 ### `failed_message_retry_strategy` Configures retry behavior for failed messages. ```rust theme={null} pub const fn failed_message_retry_strategy(mut self, strategy: RetryStrategy) -> Self ``` ### `enable_scheduling` Enables message scheduling. ```rust theme={null} pub const fn enable_scheduling(mut self, enable: bool) -> Self ``` **Default:** false ### `build` Builds the queue and connects to the broker. ```rust theme={null} pub async fn build(self) -> Result ``` # Broker Overview Source: https://densumesh-broccoli-27.localhost:3000/brokers/overview Supported message brokers and how to choose ## Supported brokers Broccoli supports multiple message brokers through feature flags: | Broker | Feature | URL Scheme | Status | | --------- | ----------- | ----------------- | ---------- | | Redis | `redis` | `redis://` | ✅ Stable | | RabbitMQ | `rabbitmq` | `amqp://` | ✅ Stable | | SurrealDB | `surrealdb` | `ws://`, `mem://` | ✅ Stable | | Kafka | - | - | 🚧 Planned | ## How broker selection works Broccoli automatically selects the broker based on your connection URL: ```rust theme={null} // Redis let queue = BroccoliQueue::builder("redis://localhost:6379").build().await?; // RabbitMQ let queue = BroccoliQueue::builder("amqp://localhost:5672").build().await?; // SurrealDB let queue = BroccoliQueue::builder("ws://localhost:8000").build().await?; ``` You must enable the appropriate feature flag for your broker. If you use a URL scheme without the corresponding feature enabled, you'll get a compile-time error. ## Choosing a broker **Best for:** General purpose, high throughput, simple setup * Fast in-memory storage * Simple operational model * Good for most use cases * Supports fairness queues **Best for:** Complex routing, enterprise messaging * Robust message delivery guarantees * Advanced routing capabilities * Excellent monitoring tools * Supports delayed/scheduled messages (with plugin) **Best for:** When you already use SurrealDB, embedded queues * Use your existing database * In-memory option for testing * Good for embedded scenarios ## Feature comparison | Feature | Redis | RabbitMQ | SurrealDB | | ------------------ | ----- | --------------- | --------- | | Connection pooling | ✅ | ✅ | ✅ | | Retry handling | ✅ | ✅ | ✅ | | Message scheduling | ✅ | ✅ (with plugin) | ✅ | | Fairness queues | ✅ | ❌ | ❌ | | Management API | ✅ | ✅ | ❌ | | In-memory mode | ❌ | ❌ | ✅ | ## Enabling multiple brokers You can enable multiple brokers in your project: ```toml theme={null} [dependencies] broccoli_queue = { version = "0.4", features = ["redis", "rabbitmq"] } ``` This allows your application to connect to different brokers at runtime: ```rust theme={null} async fn connect_to_broker(url: &str) -> Result { BroccoliQueue::builder(url) .pool_connections(5) .build() .await } // Connect based on environment let broker_url = std::env::var("BROKER_URL")?; let queue = connect_to_broker(&broker_url).await?; ``` ## Connection URL formats ```text Redis theme={null} redis://localhost:6379 redis://user:password@localhost:6379 redis://localhost:6379/0 # database number rediss://localhost:6379 # TLS ``` ```text RabbitMQ theme={null} amqp://localhost:5672 amqp://user:password@localhost:5672 amqp://user:password@localhost:5672/vhost amqps://localhost:5671 # TLS ``` ```text SurrealDB theme={null} ws://localhost:8000 wss://localhost:8000 # TLS mem:// # in-memory (testing) ``` ## Broker configuration All brokers share common configuration through the builder: ```rust theme={null} let queue = BroccoliQueue::builder("redis://localhost:6379") // Connection pool size .pool_connections(10) // Retry configuration .failed_message_retry_strategy( RetryStrategy::new() .with_attempts(5) .retry_failed(true) ) // Enable scheduling (broker-specific behavior) .enable_scheduling(true) .build() .await?; ``` See the individual broker pages for broker-specific options and considerations. # RabbitMQ Source: https://densumesh-broccoli-27.localhost:3000/brokers/rabbitmq Using RabbitMQ as your message broker ## Overview RabbitMQ is a robust message broker with advanced routing capabilities, making it ideal for complex messaging scenarios. ## Installation Enable RabbitMQ support in your `Cargo.toml`: ```toml theme={null} [dependencies] broccoli_queue = { version = "0.4", default-features = false, features = ["rabbitmq"] } ``` Or alongside Redis: ```toml theme={null} [dependencies] broccoli_queue = { version = "0.4", features = ["rabbitmq"] } ``` ## Connection ```rust theme={null} use broccoli_queue::queue::BroccoliQueue; let queue = BroccoliQueue::builder("amqp://localhost:5672") .pool_connections(10) .build() .await?; ``` ### Connection URL formats ```text theme={null} amqp://localhost:5672 # Basic amqp://user:password@localhost:5672 # With credentials amqp://user:password@localhost:5672/vhost # With virtual host amqps://localhost:5671 # TLS connection ``` ## Starting RabbitMQ ### Docker ```bash theme={null} docker run -d --name rabbitmq \ -p 5672:5672 \ -p 15672:15672 \ rabbitmq:management ``` Access the management UI at `http://localhost:15672` (guest/guest). ### With custom credentials ```bash theme={null} docker run -d --name rabbitmq \ -p 5672:5672 \ -p 15672:15672 \ -e RABBITMQ_DEFAULT_USER=myuser \ -e RABBITMQ_DEFAULT_PASS=mypassword \ rabbitmq:management ``` ## Features ### Connection pooling RabbitMQ uses `deadpool-lapin` for connection pooling: ```rust theme={null} let queue = BroccoliQueue::builder("amqp://localhost:5672") .pool_connections(10) .build() .await?; ``` ### Message scheduling Message scheduling with RabbitMQ requires the **delayed-exchange plugin**. #### Install the plugin ```bash theme={null} # In Docker docker exec rabbitmq rabbitmq-plugins enable rabbitmq_delayed_message_exchange # Or on host rabbitmq-plugins enable rabbitmq_delayed_message_exchange ``` See the [RabbitMQ scheduling guide](https://www.rabbitmq.com/blog/2015/04/16/scheduling-messages-with-rabbitmq) for details. #### Use scheduling ```rust theme={null} use broccoli_queue::queue::PublishOptions; use time::Duration; let queue = BroccoliQueue::builder("amqp://localhost:5672") .enable_scheduling(true) // Must be enabled .build() .await?; let options = PublishOptions::builder() .delay(Duration::minutes(5)) .build(); queue.publish("jobs", None, &job, Some(options)).await?; ``` ## RabbitMQ concepts Broccoli abstracts RabbitMQ concepts, but understanding them helps with debugging: | Broccoli | RabbitMQ | | --------------- | ----------------------------------------------- | | Queue name | Queue name | | `publish()` | Publish to default exchange | | `consume()` | Basic consume with prefetch | | `acknowledge()` | Basic ack | | `reject()` | Basic nack (with requeue based on retry config) | ## Management API With the `management` feature: ```rust theme={null} #[cfg(feature = "management")] { let status = queue.queue_status("jobs".into(), None).await?; println!("Queue status: {:?}", status); } ``` ## Configuration example ```rust theme={null} use broccoli_queue::queue::{BroccoliQueue, RetryStrategy}; let queue = BroccoliQueue::builder("amqp://user:pass@localhost:5672/myapp") .pool_connections(15) .failed_message_retry_strategy( RetryStrategy::new() .with_attempts(5) .retry_failed(true) ) .enable_scheduling(true) .build() .await?; ``` ## Best practices Separate environments using virtual hosts: ```rust theme={null} // Development BroccoliQueue::builder("amqp://localhost:5672/dev") // Production BroccoliQueue::builder("amqp://localhost:5672/prod") ``` RabbitMQ provides publisher confirms for guaranteed delivery. Broccoli uses these internally for reliability. Use the RabbitMQ management UI (`localhost:15672`) to: * View queue depths * Monitor connection counts * Check message rates * Debug delivery issues Set RabbitMQ memory limits to prevent OOM: ```bash theme={null} docker run -d --name rabbitmq \ -e RABBITMQ_VM_MEMORY_HIGH_WATERMARK=0.6 \ rabbitmq:management ``` ## Troubleshooting ### Connection refused ``` Connection refused (os error 111) ``` Verify RabbitMQ is running: ```bash theme={null} docker ps | grep rabbitmq # or rabbitmqctl status ``` ### ACCESS\_REFUSED ``` ACCESS_REFUSED - Login was refused ``` Check credentials in your connection URL: ```rust theme={null} BroccoliQueue::builder("amqp://user:password@localhost:5672") ``` ### NOT\_FOUND for scheduled messages ``` NOT_FOUND - no exchange 'delayed' ``` Enable the delayed message exchange plugin: ```bash theme={null} rabbitmq-plugins enable rabbitmq_delayed_message_exchange ``` ### High memory usage RabbitMQ can accumulate messages in memory. Monitor and set limits: ```bash theme={null} # Check status rabbitmqctl status # Set watermark (fraction of available RAM) rabbitmqctl set_vm_memory_high_watermark 0.5 ``` # Redis Source: https://densumesh-broccoli-27.localhost:3000/brokers/redis Using Redis as your message broker ## Overview Redis is the default broker for Broccoli. It provides fast, in-memory message storage with persistence options. ## Installation Redis support is enabled by default: ```toml theme={null} [dependencies] broccoli_queue = "0.4" ``` Or explicitly: ```toml theme={null} [dependencies] broccoli_queue = { version = "0.4", features = ["redis"] } ``` ## Connection ```rust theme={null} use broccoli_queue::queue::BroccoliQueue; let queue = BroccoliQueue::builder("redis://localhost:6379") .pool_connections(10) .build() .await?; ``` ### Connection URL formats ```text theme={null} redis://localhost:6379 # Basic redis://user:password@localhost:6379 # With auth redis://localhost:6379/2 # Database number rediss://localhost:6379 # TLS connection ``` ## Starting Redis ### Docker ```bash theme={null} docker run -d --name redis -p 6379:6379 redis:latest ``` ### With persistence ```bash theme={null} docker run -d --name redis -p 6379:6379 \ -v redis-data:/data \ redis:latest redis-server --appendonly yes ``` ### With password ```bash theme={null} docker run -d --name redis -p 6379:6379 \ redis:latest redis-server --requirepass yourpassword ``` ## Features ### Connection pooling Redis uses `bb8-redis` for connection pooling: ```rust theme={null} let queue = BroccoliQueue::builder("redis://localhost:6379") .pool_connections(20) // 20 connections in pool .build() .await?; ``` ### Fairness queues Redis supports fairness queues through disambiguators, ensuring fair processing across different message categories: ```rust theme={null} // Publish with disambiguator queue.publish("jobs", Some("tenant-a".into()), &job_a, None).await?; queue.publish("jobs", Some("tenant-b".into()), &job_b, None).await?; // Consume with fairness enabled let options = ConsumeOptions::builder() .fairness(true) .build(); queue.process_messages("jobs", Some(4), Some(options), |msg| async move { // Messages are distributed fairly across tenant-a and tenant-b Ok(()) }).await?; ``` ### Message scheduling Redis supports delayed message delivery: ```rust theme={null} use broccoli_queue::queue::PublishOptions; use time::Duration; let queue = BroccoliQueue::builder("redis://localhost:6379") .enable_scheduling(true) .build() .await?; // Delay message by 30 seconds let options = PublishOptions::builder() .delay(Duration::seconds(30)) .build(); queue.publish("jobs", None, &job, Some(options)).await?; ``` ## Redis data structures Broccoli uses the following Redis keys: | Key pattern | Type | Purpose | | ---------------------------------- | ---------- | ------------------------ | | `{queue}` | List | Main queue | | `{queue}:processing` | Sorted Set | Messages being processed | | `{queue}:failed` | List | Failed messages | | `{queue}:scheduled` | Sorted Set | Scheduled messages | | `{queue}:fairness:{disambiguator}` | List | Fairness sub-queues | ## Management API With the `management` feature enabled: ```rust theme={null} #[cfg(feature = "management")] { let status = queue.queue_status("jobs".into(), None).await?; println!("Pending: {}", status.pending_count); println!("Processing: {}", status.processing_count); println!("Failed: {}", status.failed_count); } ``` ## Best practices Set pool connections based on your workload: * Low traffic: 5-10 connections * Medium traffic: 10-20 connections * High traffic: 20-50 connections ```rust theme={null} .pool_connections(20) ``` Use Redis persistence (RDB or AOF) in production to prevent message loss: ```bash theme={null} redis-server --appendonly yes ``` Monitor Redis memory, especially with large message payloads: ```bash theme={null} redis-cli INFO memory ``` For high-volume applications, consider Redis Cluster for horizontal scaling. ## Troubleshooting ### Connection refused ``` Connection refused (os error 111) ``` Ensure Redis is running and accessible: ```bash theme={null} redis-cli ping ``` ### Authentication failed ``` NOAUTH Authentication required ``` Include credentials in your connection URL: ```rust theme={null} BroccoliQueue::builder("redis://user:password@localhost:6379") ``` ### Pool exhausted If you see timeout errors, increase pool size: ```rust theme={null} .pool_connections(30) ``` # SurrealDB Source: https://densumesh-broccoli-27.localhost:3000/brokers/surrealdb Using SurrealDB as your message broker ## Overview SurrealDB can be used as a message broker, which is useful when you already use SurrealDB in your application or need an embedded queue for testing. ## Installation Enable SurrealDB support in your `Cargo.toml`: ```toml theme={null} [dependencies] broccoli_queue = { version = "0.4", default-features = false, features = ["surrealdb"] } ``` ## Connection ### WebSocket connection ```rust theme={null} use broccoli_queue::queue::BroccoliQueue; let queue = BroccoliQueue::builder("ws://localhost:8000") .pool_connections(5) .build() .await?; ``` ### In-memory (testing) ```rust theme={null} let queue = BroccoliQueue::builder("mem://") .build() .await?; ``` ### Reuse existing connection If you already have a SurrealDB connection, you can reuse it: ```rust theme={null} use surrealdb::Surreal; use surrealdb::engine::any::Any; // Your existing SurrealDB connection let db: Surreal = // ... let queue = BroccoliQueue::builder_with(db) .failed_message_retry_strategy(Default::default()) .build() .await?; ``` ## Starting SurrealDB ### Docker ```bash theme={null} docker run -d --name surrealdb \ -p 8000:8000 \ surrealdb/surrealdb:latest \ start --user root --pass root ``` ### With file persistence ```bash theme={null} docker run -d --name surrealdb \ -p 8000:8000 \ -v surrealdb-data:/data \ surrealdb/surrealdb:latest \ start --user root --pass root file:/data/mydb.db ``` ## Features ### In-memory mode Perfect for testing and development: ```rust theme={null} #[tokio::test] async fn test_message_processing() { let queue = BroccoliQueue::builder("mem://") .build() .await .unwrap(); // Run your tests with an ephemeral queue queue.publish("test", None, &"test message", None).await.unwrap(); } ``` ### Message scheduling SurrealDB supports delayed message delivery: ```rust theme={null} use broccoli_queue::queue::PublishOptions; use time::Duration; let queue = BroccoliQueue::builder("ws://localhost:8000") .enable_scheduling(true) .build() .await?; let options = PublishOptions::builder() .delay(Duration::seconds(60)) .build(); queue.publish("jobs", None, &job, Some(options)).await?; ``` ## SurrealDB data model Broccoli stores messages in SurrealDB tables: | Table | Purpose | | ----------------------------- | ------------------------ | | `broccoli_{queue}` | Main queue messages | | `broccoli_{queue}_processing` | Messages being processed | | `broccoli_{queue}_failed` | Failed messages | | `broccoli_{queue}_scheduled` | Scheduled messages | ## Configuration example ```rust theme={null} use broccoli_queue::queue::{BroccoliQueue, RetryStrategy}; let queue = BroccoliQueue::builder("ws://localhost:8000") .pool_connections(5) .failed_message_retry_strategy( RetryStrategy::new() .with_attempts(3) ) .enable_scheduling(true) .build() .await?; ``` ## Best practices The `mem://` scheme is ideal for unit and integration tests: ```rust theme={null} #[cfg(test)] mod tests { use super::*; async fn setup_test_queue() -> BroccoliQueue { BroccoliQueue::builder("mem://") .build() .await .unwrap() } } ``` If your application already uses SurrealDB, reuse the connection: ```rust theme={null} let queue = BroccoliQueue::builder_with(existing_db) .build() .await?; ``` SurrealDB may have different concurrency characteristics than dedicated message brokers. Test with your expected load. ## Limitations SurrealDB as a message broker has some limitations compared to Redis or RabbitMQ: * **No fairness queues** - Disambiguator-based fairness is not supported * **No management API** - Queue status monitoring is not available * **Performance** - May not match dedicated brokers for high throughput ## When to use SurrealDB Good use cases: * You already use SurrealDB and want to minimize dependencies * Testing and development with in-memory queues * Embedded applications with moderate queue requirements Consider Redis or RabbitMQ for: * High-throughput production workloads * Complex routing requirements * When you need fairness queues or management APIs ## Troubleshooting ### Connection failed ``` Failed to connect to broker ``` Verify SurrealDB is running: ```bash theme={null} curl http://localhost:8000/health ``` ### Concurrency errors ``` Broker error: non-idempotent operation ``` This can occur with parallel consumers. Consider reducing concurrency: ```rust theme={null} queue.process_messages("jobs", Some(2), None, handler).await?; ``` # Messages Source: https://densumesh-broccoli-27.localhost:3000/concepts/messages Working with BrokerMessage and message payloads ## Overview Messages in Broccoli are wrapped in a `BrokerMessage` struct that includes metadata alongside your custom payload. ## BrokerMessage structure ```rust theme={null} pub struct BrokerMessage { /// Unique identifier for the message pub task_id: uuid::Uuid, /// Your custom message content pub payload: T, /// Number of processing attempts made pub attempts: u8, /// Optional disambiguator for fairness queues pub disambiguator: Option, } ``` ## Creating message payloads Your payload type must implement `Clone`, `Serialize`, and `Deserialize`: ```rust theme={null} use serde::{Serialize, Deserialize}; #[derive(Debug, Clone, Serialize, Deserialize)] struct JobPayload { id: String, task_name: String, parameters: serde_json::Value, created_at: chrono::DateTime, } ``` ## Publishing messages When you publish a message, Broccoli automatically: 1. Generates a unique `task_id` 2. Wraps your payload in a `BrokerMessage` 3. Serializes and sends to the broker ```rust theme={null} let job = JobPayload { id: "job-123".to_string(), task_name: "process_data".to_string(), parameters: serde_json::json!({"input": "data.csv"}), created_at: chrono::Utc::now(), }; // publish() returns the wrapped message with its task_id let message = queue.publish("jobs", None, &job, None).await?; println!("Published with ID: {}", message.task_id); ``` ## Consuming messages When consuming, you receive the full `BrokerMessage`: ```rust theme={null} let message: BrokerMessage = queue.consume("jobs", None).await?; // Access metadata println!("Task ID: {}", message.task_id); println!("Attempts: {}", message.attempts); // Access your payload println!("Job ID: {}", message.payload.id); println!("Task: {}", message.payload.task_name); ``` ## Message acknowledgment After processing a message, you must acknowledge or reject it: ### Acknowledge (success) Removes the message from the processing queue: ```rust theme={null} queue.acknowledge("jobs", message).await?; ``` ### Reject (failure) Moves the message to the retry queue or failed queue: ```rust theme={null} queue.reject("jobs", message).await?; ``` The message will be: * **Re-queued** if `attempts < max_retries` * **Moved to failed queue** if `attempts >= max_retries` ### Cancel Remove a message by ID without processing: ```rust theme={null} queue.cancel("jobs", message_id.to_string()).await?; ``` ## Message ordering Broccoli does not guarantee strict message ordering. Messages may be processed out of order, especially with multiple workers. If you need ordered processing: * Use a single worker (`concurrency: Some(1)`) * Or implement ordering logic in your application ## Complex payload example ```rust theme={null} use chrono::{DateTime, Utc}; use serde::{Serialize, Deserialize}; #[derive(Debug, Clone, Serialize, Deserialize)] struct EmailJob { recipient: String, subject: String, body: String, attachments: Vec, priority: Priority, scheduled_for: Option>, } #[derive(Debug, Clone, Serialize, Deserialize)] struct Attachment { filename: String, content_type: String, data: Vec, } #[derive(Debug, Clone, Serialize, Deserialize)] enum Priority { Low, Normal, High, Urgent, } ``` # Queue Source: https://densumesh-broccoli-27.localhost:3000/concepts/queue Understanding the BroccoliQueue and its operations ## Overview The `BroccoliQueue` is the main entry point for interacting with Broccoli. It provides methods for publishing messages, consuming messages, and processing messages with handlers. ## Creating a queue Use the builder pattern to create a queue instance: ```rust theme={null} use broccoli_queue::queue::BroccoliQueue; let queue = BroccoliQueue::builder("redis://localhost:6379") .pool_connections(5) .failed_message_retry_strategy(Default::default()) .build() .await?; ``` ### Builder options | Option | Description | Default | | ----------------------------------------- | --------------------------------- | ------------------ | | `pool_connections(n)` | Number of connections in the pool | 10 | | `failed_message_retry_strategy(strategy)` | Configure retry behavior | 3 retries, enabled | | `enable_scheduling(bool)` | Enable message scheduling | false | ## Publishing messages ### Single message ```rust theme={null} use broccoli_queue::queue::BroccoliQueue; let job = JobPayload { id: "1".into(), task: "process".into() }; // Basic publish queue.publish("jobs", None, &job, None).await?; // With disambiguator for fairness queue.publish("jobs", Some("tenant-123".into()), &job, None).await?; ``` ### Batch publishing ```rust theme={null} let jobs = vec![ JobPayload { id: "1".into(), task: "task-a".into() }, JobPayload { id: "2".into(), task: "task-b".into() }, ]; queue.publish_batch("jobs", None, jobs, None).await?; ``` ## Consuming messages ### Blocking consume Waits until a message is available: ```rust theme={null} let message: BrokerMessage = queue.consume("jobs", None).await?; println!("Received: {:?}", message.payload); // Process the message... // Acknowledge success queue.acknowledge("jobs", message).await?; ``` ### Non-blocking consume Returns immediately with `None` if no message is available: ```rust theme={null} if let Some(message) = queue.try_consume::("jobs", None).await? { println!("Got message: {:?}", message.payload); queue.acknowledge("jobs", message).await?; } else { println!("No messages available"); } ``` ### Batch consume Consume multiple messages with a timeout: ```rust theme={null} use time::Duration; let messages = queue.consume_batch::( "jobs", 10, // batch size Duration::seconds(5), // timeout None, ).await?; for message in messages { // Process and acknowledge each message queue.acknowledge("jobs", message).await?; } ``` ## Processing messages The `process_messages` method provides a convenient way to consume and process messages in a loop: ```rust theme={null} queue.process_messages( "jobs", Some(4), // 4 concurrent workers None, // default consume options |message: BrokerMessage| async move { println!("Processing: {:?}", message.payload); // Return Ok(()) on success, Err on failure Ok(()) }, ).await?; ``` ### With custom handlers For more control over success and error handling: ```rust theme={null} queue.process_messages_with_handlers( "jobs", Some(4), None, // Message handler |msg| async move { process_job(msg.payload).await }, // Success handler |msg, result| async move { println!("Job {} completed", msg.task_id); Ok(()) }, // Error handler |msg, error| async move { eprintln!("Job {} failed: {:?}", msg.task_id, error); Ok(()) }, ).await?; ``` ## Message lifecycle ```mermaid theme={null} stateDiagram-v2 [*] --> Queued: publish() Queued --> Processing: consume() Processing --> Completed: acknowledge() Processing --> Failed: reject() Failed --> Queued: retry (if attempts < max) Failed --> DeadLetter: max retries exceeded Processing --> Cancelled: cancel() Completed --> [*] DeadLetter --> [*] Cancelled --> [*] ``` ## Queue size Get the current size of a queue: ```rust theme={null} let sizes = queue.size("jobs").await?; for (queue_name, size) in sizes { println!("{}: {} messages", queue_name, size); } ``` For fairness queues, this returns sizes for each disambiguator sub-queue. # Retry Strategy Source: https://densumesh-broccoli-27.localhost:3000/concepts/retry-strategy Configure how failed messages are retried ## Overview Broccoli includes built-in retry handling for failed messages. When a message handler returns an error, the message can be automatically re-queued for another attempt. ## Default behavior By default, Broccoli: * Enables retries for failed messages * Allows up to 3 retry attempts * Moves messages to a failed queue after exhausting retries ## Configuring retry strategy Use `RetryStrategy` when building your queue: ```rust theme={null} use broccoli_queue::queue::{BroccoliQueue, RetryStrategy}; let queue = BroccoliQueue::builder("redis://localhost:6379") .failed_message_retry_strategy( RetryStrategy::new() .with_attempts(5) // Allow up to 5 retries .retry_failed(true) // Enable retries ) .build() .await?; ``` ## RetryStrategy options ### `with_attempts(n)` Set the maximum number of retry attempts: ```rust theme={null} RetryStrategy::new().with_attempts(10) // Up to 10 retries ``` ### `retry_failed(bool)` Enable or disable retries entirely: ```rust theme={null} // Disable retries - failed messages go directly to failed queue RetryStrategy::new().retry_failed(false) ``` ## How retries work ```mermaid theme={null} flowchart TD A[Message consumed] --> B{Handler succeeds?} B -->|Yes| C[Acknowledge] B -->|No| D{Retries enabled?} D -->|No| E[Move to failed queue] D -->|Yes| F{attempts < max?} F -->|Yes| G[Increment attempts] G --> H[Re-queue message] F -->|No| E ``` ## Tracking attempts Each `BrokerMessage` includes an `attempts` field: ```rust theme={null} queue.process_messages("jobs", Some(1), None, |message: BrokerMessage| async move { println!("Attempt #{} for job {}", message.attempts + 1, message.task_id); if message.attempts >= 2 { // Different handling for repeated failures eprintln!("Job has failed multiple times, attempting recovery..."); } process_job(message.payload).await }).await?; ``` ## Custom error handling For fine-grained control, use `process_messages_with_handlers`: ```rust theme={null} queue.process_messages_with_handlers( "jobs", Some(4), None, // Main handler |msg| async move { risky_operation(msg.payload).await }, // Success handler |msg, result| async move { log::info!("Job {} succeeded after {} attempts", msg.task_id, msg.attempts); Ok(()) }, // Error handler - called before retry/rejection |msg, error| async move { log::error!("Job {} failed (attempt {}): {:?}", msg.task_id, msg.attempts, error); // You could send alerts, update metrics, etc. if msg.attempts >= 2 { alert_on_repeated_failure(&msg, &error).await; } Ok(()) }, ).await?; ``` ## Best practices Design your message handlers to be idempotent. Since messages may be retried, the same message could be processed multiple times. ```rust theme={null} async fn process_job(job: JobPayload) -> Result<(), BroccoliError> { // Check if already processed if is_already_processed(&job.id).await? { return Ok(()); // Skip duplicate processing } // Process the job do_work(&job).await?; // Mark as processed mark_processed(&job.id).await?; Ok(()) } ``` Consider whether failures are transient (network issues, temporary unavailability) or permanent (invalid data, business rule violations). ```rust theme={null} async fn process_job(job: JobPayload) -> Result<(), BroccoliError> { match validate_job(&job) { Err(ValidationError::InvalidData) => { // Permanent failure - don't retry log::error!("Invalid job data: {:?}", job); return Err(BroccoliError::Job("Invalid data - will not retry".into())); } Err(ValidationError::ServiceUnavailable) => { // Transient failure - let it retry return Err(BroccoliError::Job("Service unavailable".into())); } Ok(_) => {} } do_work(&job).await } ``` Broccoli doesn't include built-in backoff, but you can implement it in your handler: ```rust theme={null} async fn process_with_backoff(msg: BrokerMessage) -> Result<(), BroccoliError> { // Add delay based on attempt count let delay_ms = 100 * 2_u64.pow(msg.attempts as u32); tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await; process_job(msg.payload).await } ``` ## Failed message queue Messages that exceed the retry limit are moved to a failed queue (`{queue_name}:failed`). You can: * Monitor this queue for alerting * Manually inspect and retry failed messages * Use the management API to retrieve queue status ```rust theme={null} // With management feature enabled #[cfg(feature = "management")] { let status = queue.queue_status("jobs".into(), None).await?; println!("Failed messages: {}", status.failed_count); } ``` # Broccoli Source: https://densumesh-broccoli-27.localhost:3000/index A robust, type-safe message queue system for Rust applications # Broccoli A robust message queue system for Rust applications, designed as a Rust alternative to [Celery](https://docs.celeryq.dev/). Broccoli provides asynchronous message processing with type safety, configurable retry strategies, and support for multiple message brokers. Get up and running with Broccoli in minutes Add Broccoli to your Rust project Learn about supported message brokers Explore the complete API documentation ## Features Leverage Rust's type system for compile-time message validation Support for Redis, RabbitMQ, and SurrealDB backends Built-in retry strategies for failed message handling Efficient resource management with connection pools Schedule messages for delayed or future delivery Process messages with configurable worker concurrency ## Quick example ```rust theme={null} use broccoli_queue::{queue::BroccoliQueue, brokers::broker::BrokerMessage}; use serde::{Serialize, Deserialize}; #[derive(Debug, Clone, Serialize, Deserialize)] struct JobPayload { id: String, task_name: String, } #[tokio::main] async fn main() -> Result<(), Box> { // Initialize the queue with Redis let queue = BroccoliQueue::builder("redis://localhost:6379") .pool_connections(5) .failed_message_retry_strategy(Default::default()) .build() .await?; // Publish a message let job = JobPayload { id: "job-1".to_string(), task_name: "process_data".to_string(), }; queue.publish("jobs", None, &job, None).await?; // Process messages queue.process_messages("jobs", Some(2), None, |message: BrokerMessage| async move { println!("Processing: {:?}", message.payload); Ok(()) }).await?; Ok(()) } ``` ## Why Broccoli? If you're familiar with Celery but want to leverage Rust's performance and type safety, Broccoli provides: * **Rust performance** - Native async/await with Tokio runtime * **Type safety** - Compile-time message validation with serde * **Multiple brokers** - Choose Redis, RabbitMQ, or SurrealDB based on your needs * **Simple API** - Builder pattern for easy configuration ## Current status Broccoli is under active development. Current broker support: | Broker | Status | | --------- | ----------- | | Redis | ✅ Available | | RabbitMQ | ✅ Available | | SurrealDB | ✅ Available | | Kafka | 🚧 Planned | # Installation Source: https://densumesh-broccoli-27.localhost:3000/installation Add Broccoli to your Rust project ## Requirements * Rust 1.70 or later * A supported message broker (Redis, RabbitMQ, or SurrealDB) ## Add the dependency Add Broccoli to your `Cargo.toml`: ```toml theme={null} [dependencies] broccoli_queue = "0.4" ``` ## Feature flags Broccoli uses feature flags to enable different message brokers. By default, only Redis is enabled. ### Available features | Feature | Description | Default | | ------------ | ------------------------------- | ------- | | `redis` | Enable Redis broker support | ✅ | | `rabbitmq` | Enable RabbitMQ broker support | ❌ | | `surrealdb` | Enable SurrealDB broker support | ❌ | | `management` | Enable queue management API | ❌ | ### Using specific brokers ```toml Redis only (default) theme={null} [dependencies] broccoli_queue = "0.4" ``` ```toml RabbitMQ only theme={null} [dependencies] broccoli_queue = { version = "0.4", default-features = false, features = ["rabbitmq"] } ``` ```toml SurrealDB only theme={null} [dependencies] broccoli_queue = { version = "0.4", default-features = false, features = ["surrealdb"] } ``` ```toml Multiple brokers theme={null} [dependencies] broccoli_queue = { version = "0.4", features = ["rabbitmq", "surrealdb"] } ``` ### Management feature Enable the `management` feature to access queue status and monitoring capabilities: ```toml theme={null} [dependencies] broccoli_queue = { version = "0.4", features = ["management"] } ``` ## Required dependencies Broccoli requires async runtime support. Add Tokio and serde to your project: ```toml theme={null} [dependencies] broccoli_queue = "0.4" tokio = { version = "1", features = ["full"] } serde = { version = "1.0", features = ["derive"] } ``` For working with timestamps in your message payloads, you may also want: ```toml theme={null} [dependencies] chrono = { version = "0.4", features = ["serde"] } ``` ## Verifying installation Create a simple test to verify Broccoli is installed correctly: ```rust theme={null} use broccoli_queue::queue::BroccoliQueue; #[tokio::main] async fn main() { // This will fail to connect without a running broker, // but verifies the library is installed correctly let result = BroccoliQueue::builder("redis://localhost:6379") .build() .await; match result { Ok(_) => println!("Connected successfully!"), Err(e) => println!("Connection failed (expected without broker): {}", e), } } ``` Run with: ```bash theme={null} cargo run ``` ## Next steps Build your first producer and consumer Configure your message broker # Quickstart Source: https://densumesh-broccoli-27.localhost:3000/quickstart Get up and running with Broccoli in 5 minutes ## Prerequisites Before you begin, make sure you have: * Rust installed (1.70+ recommended) * A running message broker (Redis, RabbitMQ, or SurrealDB) ## Install Broccoli Add Broccoli to your `Cargo.toml`: ```toml theme={null} [dependencies] broccoli_queue = "0.4" tokio = { version = "1", features = ["full"] } serde = { version = "1.0", features = ["derive"] } ``` By default, Broccoli uses Redis. To use a different broker, specify features: ```toml Redis (default) theme={null} [dependencies] broccoli_queue = "0.4" ``` ```toml RabbitMQ theme={null} [dependencies] broccoli_queue = { version = "0.4", default-features = false, features = ["rabbitmq"] } ``` ```toml SurrealDB theme={null} [dependencies] broccoli_queue = { version = "0.4", default-features = false, features = ["surrealdb"] } ``` ## Define your message type Create a struct that represents your job payload. It must implement `Serialize` and `Deserialize`: ```rust theme={null} use serde::{Serialize, Deserialize}; #[derive(Debug, Clone, Serialize, Deserialize)] struct JobPayload { id: String, task_name: String, } ``` ## Create a producer Publish messages to a queue: ```rust theme={null} use broccoli_queue::queue::BroccoliQueue; #[tokio::main] async fn main() -> Result<(), Box> { // Connect to Redis let queue = BroccoliQueue::builder("redis://localhost:6379") .pool_connections(5) .build() .await?; // Create a job let job = JobPayload { id: "job-1".to_string(), task_name: "process_data".to_string(), }; // Publish to the "jobs" queue queue.publish("jobs", None, &job, None).await?; println!("Job published!"); Ok(()) } ``` ## Create a consumer Process messages from the queue: ```rust theme={null} use broccoli_queue::{queue::BroccoliQueue, brokers::broker::BrokerMessage}; #[tokio::main] async fn main() -> Result<(), Box> { // Connect to Redis let queue = BroccoliQueue::builder("redis://localhost:6379") .pool_connections(5) .failed_message_retry_strategy(Default::default()) .build() .await?; // Process messages with 2 concurrent workers queue.process_messages( "jobs", // Queue name Some(2), // Number of workers None, // Consume options |message: BrokerMessage| async move { println!("Processing job: {:?}", message.payload); // Your processing logic here Ok(()) }, ).await?; Ok(()) } ``` ## Run your application 1. Start your message broker (e.g., Redis): ```bash theme={null} docker run -d -p 6379:6379 redis ``` 2. Run the consumer: ```bash theme={null} cargo run --bin consumer ``` 3. Run the producer: ```bash theme={null} cargo run --bin producer ``` ## Next steps Learn about queues, messages, and the processing model Configure how failed messages are retried Choose and configure your message broker Schedule messages for delayed delivery