# How Background Jobs Are Processed Using SQS Workers in the Macro Repository

> Discover how Macro Inc's SQS workers process background jobs. Learn about concurrent message processing, automatic restarts, and graceful shutdown in macro-inc/macro.

- Repository: [Macro/macro](https://github.com/macro-inc/macro)
- Tags: internals
- Published: 2026-08-17

---

**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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/services/search_processing_service/src/main.rs), the worker is instantiated with configuration values for maximum messages per request and long-poll wait time:

```rust
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`](https://github.com/macro-inc/macro/blob/main/services/email_service/src/pubsub/inbox_sync/worker.rs), the processing loop demonstrates how domain-specific logic integrates with the worker:

```rust
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`](https://github.com/macro-inc/macro/blob/main/crates/sqs_worker/src/lib.rs) extracts the receipt handle and calls the SQS `DeleteMessage` API:

```rust
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`](https://github.com/macro-inc/macro/blob/main/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:

```rust
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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/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 like `process_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`](https://github.com/macro-inc/macro/blob/main/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`](https://github.com/macro-inc/macro/blob/main/services/email_service/src/pubsub/inbox_sync/worker.rs)) and SFS upload processing ([`services/email_service/src/pubsub/sfs_uploader/worker.rs`](https://github.com/macro-inc/macro/blob/main/services/email_service/src/pubsub/sfs_uploader/worker.rs)). The `worker_trigger` service also demonstrates how background workers can be initiated from external events.