How Background Jobs Are Processed Using SQS Workers in the Macro Repository
The Macro codebase processes asynchronous background jobs through a lightweight SQSWorker abstraction that polls AWS SQS queues, processes messages concurrently, and handles deletion and error recovery with automatic restarts and graceful shutdown support.
The macro-inc/macro repository implements a standardized pattern for background job processing using AWS Simple Queue Service (SQS). At the heart of this system lies the SQSWorker component, which encapsulates queue polling, message handling, and lifecycle management across multiple services including search processing and email synchronization.
The SQSWorker Abstraction
The core implementation resides in crates/sqs_worker/src/lib.rs, where the SQSWorker struct wraps the AWS SDK client with service-specific configuration. This abstraction handles the complexity of SQS API interactions while providing a consistent interface for consuming messages.
Worker Construction
Each background service constructs an SQSWorker by supplying the aws_sdk_sqs::Client, target queue URL, and polling parameters. In services/search_processing_service/src/main.rs, the worker is instantiated with configuration values for maximum messages per request and long-poll wait time:
let worker = sqs_worker::SQSWorker::new(
aws_sdk_sqs::Client::new(&aws_config),
search_event_queue.to_string(),
config.queue_max_messages,
config.queue_wait_time_seconds,
);
This configuration determines how many messages the worker retrieves per API call and how long it waits for messages to arrive before returning an empty response.
The Polling Loop
Background job consumption centers on an asynchronous polling loop that continuously invokes worker.receive_messages(). This method calls the SQS ReceiveMessage API using the parameters established during construction.
The polling mechanism integrates with Tokio's cancellation primitives through run_until_cancelled, which monitors a tokio_util::sync::CancellationToken. This pattern allows the worker to block efficiently on incoming messages while remaining responsive to shutdown signals.
Message Processing Pipeline
Once the worker retrieves a batch of aws_sdk_sqs::types::Message objects, it distributes them through a concurrent processing stream. The implementation uses futures::stream::iter combined with .then() to spawn asynchronous handlers for each message.
Concurrent Processing
In services/email_service/src/pubsub/inbox_sync/worker.rs, the processing loop demonstrates how domain-specific logic integrates with the worker:
let results = futures::stream::iter(messages.iter())
.then(|msg| async {
let res = process_message(ctx.clone(), msg).await;
match res {
Ok(_) => worker.cleanup_message(msg).await,
Err(e) => Err((msg.message_id.clone().unwrap_or_default(), e)),
}
})
.collect::<Vec<_>>()
.await;
This pattern enables concurrent processing of multiple messages while preserving individual error contexts. The process_message function performs domain-specific work—such as indexing documents or synchronizing email inboxes—before the worker handles message lifecycle management.
Error Handling and Message Lifecycle
The worker implements at-least-once delivery semantics through explicit acknowledgment. Successful processing triggers message deletion, while failures preserve the message in the queue for redelivery according to the queue's retry policy.
Acknowledgment via cleanup_message
The cleanup_message method in crates/sqs_worker/src/lib.rs extracts the receipt handle and calls the SQS DeleteMessage API:
pub async fn cleanup_message(&self, message: &aws_sdk_sqs::types::Message) -> anyhow::Result<()> {
if let Some(handle) = &message.receipt_handle {
delete_message::delete_message(&self.inner, &self.queue_url, handle).await?;
} else {
anyhow::bail!("no receipt handle found for message");
}
Ok(())
}
If processing fails, the worker logs the error with the message ID but skips deletion, allowing SQS to retry the message after the visibility timeout expires.
Resilience and Restart Logic
Workers implement crash resilience through outer supervision loops. As shown in services/email_service/src/pubsub/inbox_sync/worker.rs, the polling loop runs inside a loop construct that catches panics, logs failures, waits several seconds, and restarts the worker:
loop {
match run_worker().await {
Ok(_) => break,
Err(e) => {
tracing::error!(error = ?e, "worker crashed, restarting in 5 seconds");
tokio::time::sleep(Duration::from_secs(5)).await;
}
}
}
This pattern ensures transient network errors or temporary service unavailability do not permanently halt background job processing.
Graceful Shutdown
All workers accept a CancellationToken during initialization. When the application receives a shutdown signal (such as SIGTERM), the token triggers cancellation, causing run_until_cancelled to return None and break the polling loop. This mechanism allows in-flight messages to complete processing before the application exits, preventing partial work and orphaned messages.
Summary
- SQSWorker provides a unified abstraction over AWS SQS operations, living in
crates/sqs_worker/src/lib.rs. - Services construct workers with queue URLs and polling parameters, as demonstrated in
services/search_processing_service/src/main.rs. - The polling loop uses long-polling via
receive_messages()and respects cancellation tokens for graceful shutdown. - Messages process concurrently using
futures::stream, with domain logic implemented in service-specific handlers likeprocess_message. - Successful processing invokes
cleanup_message()to delete messages from the queue; failures leave messages for automatic retry. - Outer supervision loops automatically restart workers after crashes, ensuring high availability of background job processing.
Frequently Asked Questions
How does SQSWorker handle message deletion after successful processing?
The SQSWorker calls cleanup_message(), which extracts the receipt handle from the aws_sdk_sqs::types::Message and invokes the delete_message API via the internal SQS client. This occurs only after the service-specific processor returns successfully. If deletion fails or the receipt handle is missing, the method returns an error that the calling code can log without crashing the worker.
What happens when an SQS worker encounters a processing error?
When process_message returns an error, the worker logs the failure along with the message ID but does not call cleanup_message(). This leaves the message in the SQS queue, where it becomes visible again after the visibility timeout expires. SQS automatically retries the message according to the queue's redrive policy, eventually moving it to a dead-letter queue if processing continues to fail.
How does Macro ensure graceful shutdown of background workers?
Each worker accepts a tokio_util::sync::CancellationToken during construction. The polling loop wraps the receive_messages() call in run_until_cancelled(), which returns None when the token signals cancellation. This breaks the processing loop cleanly after any in-flight messages complete, allowing the application to exit without interrupting active work or losing messages.
Which services in the Macro repository use SQS workers?
The Search Processing Service (services/search_processing_service/src/main.rs) uses workers for indexing operations, while the Email Service implements multiple workers including inbox synchronization (services/email_service/src/pubsub/inbox_sync/worker.rs) and SFS upload processing (services/email_service/src/pubsub/sfs_uploader/worker.rs). The worker_trigger service also demonstrates how background workers can be initiated from external events.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →