# How Macro’s Calendar Outbox Handles Scheduling and Sending Calendar Invites

> Learn how Macro's calendar outbox schedules and sends invites using a durable pipeline with row-level locking and batch processing for at-least-once delivery.

- Repository: [Macro/macro](https://github.com/macro-inc/macro)
- Tags: how-to-guide
- Published: 2026-08-18

---

**Macro’s calendar outbox implements a durable, back-pressure-aware pipeline that converts calendar backfill jobs into SQS messages using row-level locking and batch processing to guarantee at-least-once delivery.**

The [macro-inc/macro](https://github.com/macro-inc/macro) repository uses a robust outbox pattern to reliably schedule and send calendar invites. By decoupling job creation from message publication through three dedicated outbox tables, the system ensures that calendar synchronization operations remain durable, idempotent, and resilient to transient failures.

## The Continuous Drain Loop in [`calendar_outbox.rs`](https://github.com/macro-inc/macro/blob/main/calendar_outbox.rs)

The core orchestration logic resides in [services/email_service/src/calendar_outbox.rs](https://github.com/macro-inc/macro/blob/main/services/email_service/src/calendar_outbox.rs). The `run` function enters an infinite loop that continuously drains three distinct outbox tables:

- `calendar_sync_outbox` – queues calendar-specific backfill operations.
- `email_backfill_init_outbox` – initiates email thread listing for new backfill jobs.
- `email_backfill_completion_outbox` – signals finalization when a backfill job completes.

Each iteration processes up to **50 rows** (the `BATCH_SIZE` constant), ensuring the outbox does not overwhelm downstream consumers. Before each drain, the loop checks a **cancellation token** to allow graceful shutdown.

### Batch Processing with Row-Level Locking

To prevent duplicate processing across multiple worker instances, the drain query acquires a **row-level lock** using `FOR UPDATE OF … SKIP LOCKED`. This PostgreSQL-specific clause ensures that concurrently running workers skip rows already locked by another process.

After successfully publishing a message to SQS, the worker updates the row’s `published_at` timestamp to the current UTC time. This single write operation guarantees **at-least-once delivery** while making the operation idempotent; consumers can safely process the same logical event multiple times without side effects.

## Outbox Row Mapping and Message Types

The outbox differentiates job types through strongly-typed row structures and discriminated `kind` fields.

### Calendar Sync Rows

The `OutboxRow` struct represents entries in the `calendar_sync_outbox` table. Its `kind` field determines the target backfill operation. When the kind equals `"google_calendar"`, the row transforms into a `BackfillOperation::CalendarGoogleBackfill` payload via the `to_queue_message` conversion method.

### Email Backfill Coordination

Two additional row types manage the email backfill lifecycle:

- `EmailInitOutboxRow` – triggers the initial thread enumeration for a backfill job.
- `EmailCompletionOutboxRow` – emits a finalization message once all threads are processed.

Both follow the same locking and timestamping pattern as the calendar outbox, ensuring consistency across the entire backfill pipeline.

## Scheduler Integration for Job Creation

The outbox does not create jobs directly; instead, it coordinates with the **Google Calendar Sync Scheduler**. When the global `calendar_sync_enabled` boolean flag is true, each loop iteration invokes `GoogleCalendarSyncScheduler::run_once`, passing the current UTC time.

This scheduler, defined in [crates/calendar_events/src/domain/service/google_calendar_sync_scheduler.rs](https://github.com/macro-inc/macro/blob/main/crates/calendar_events/src/domain/service/google_calendar_sync_scheduler.rs), queries the `calendar_sync_jobs` table to identify due synchronizations. It inserts new rows into `calendar_sync_outbox` for each eligible job, which the subsequent drain phase then publishes. The repository interface implemented by the scheduler lives in [crates/calendar_events/src/domain/ports/google_calendar_sync_repository.rs](https://github.com/macro-inc/macro/blob/main/crates/calendar_events/src/domain/ports/google_calendar_sync_repository.rs).

## Recovery and Idempotency Guarantees

The Macro calendar outbox is designed for **graceful degradation** and manual recovery.

### Handling Malformed Entries

If the drain encounters a row with an unrecognized `kind` value, it does not panic or retry indefinitely. Instead, the worker immediately marks the row as published by setting `published_at`. This prevents a single malformed entry from blocking the entire pipeline while logging the anomaly for operator review.

### Manual Republishing

When calendar sync is temporarily disabled or a transient error occurs, operators can requeue specific jobs using the `republish_calendar_job` helper. This function resets the `published_at` timestamp to `NULL` for a given `backfill_job_id`, allowing the standard drain loop to pick up and reprocess the job on the next iteration.

## Database Schema and Migration Patterns

The outbox tables follow a simple “log-style” schema consisting of an `id` primary key, a `backfill_job_id` foreign key, and a nullable `published_at` timestamp. Job-specific metadata columns store additional context required for message construction.

While the calendar outbox schema is embedded in the application code, a representative migration for a similar outbox pattern exists in [crates/macro_db_client/migrations/20260429140000_contacts_backfill_outbox.sql](https://github.com/macro-inc/macro/blob/main/crates/macro_db_client/migrations/20260429140000_contacts_backfill_outbox.sql). This file demonstrates the typical table layout reused across calendar, email-init, and email-completion outboxes.

### Running the Outbox Daemon

To start the calendar outbox process in your own deployment, invoke the `run` function with a database connection pool, SQS client, scheduler instance, and cancellation token:

```rust
use macro_db_client::PgPool;
use sqs_client::SQS;
use calendar_events::domain::ports::GoogleCalendarSyncRepository;
use calendar_events::domain::service::GoogleCalendarSyncScheduler;
use tokio_util::sync::CancellationToken;

#[tokio::main]
async fn main() {
    // Assume `db`, `sqs`, `scheduler` are constructed elsewhere.
    let db: PgPool = /* … */;
    let sqs: SQS = /* … */;
    let scheduler: GoogleCalendarSyncScheduler<impl GoogleCalendarSyncRepository> = /* … */;
    let token = CancellationToken::new();

    // Enable calendar sync globally (could be driven by a feature flag)
    let calendar_sync_enabled = true;

    // Start the outbox loop – it will run until the token is cancelled.
    macro_service::email_service::calendar_outbox::run(
        db,
        sqs,
        scheduler,
        calendar_sync_enabled,
        token.clone(),
    )
    .await;
}

```

### Manually Republishing a Calendar Job

After temporarily disabling calendar sync or recovering from a downstream failure, reset a specific job for reprocessing:

```rust
use macro_db_client::PgPool;
use uuid::Uuid;

async fn retry_job(db: &PgPool, job_id: Uuid) -> anyhow::Result<()> {
    macro_service::email_service::calendar_outbox::republish_calendar_job(db, job_id).await
}

```

## Summary

- **The Macro calendar outbox** uses a continuous drain loop in [`calendar_outbox.rs`](https://github.com/macro-inc/macro/blob/main/calendar_outbox.rs) to process three interrelated tables: `calendar_sync_outbox`, `email_backfill_init_outbox`, and `email_backfill_completion_outbox`.
- **Row-level locking** via `FOR UPDATE SKIP LOCKED` and a `published_at` timestamp ensure idempotent, at-least-once delivery without duplicate processing.
- **Batch sizing** is fixed at 50 rows per iteration to provide natural back-pressure against the SQS queue.
- **Scheduler integration** with `GoogleCalendarSyncScheduler` creates outbox rows for due calendar synchronizations only when the global `calendar_sync_enabled` flag is active.
- **Recovery helpers** like `republish_calendar_job` allow operators to manually requeue jobs after transient failures or configuration changes.

## Frequently Asked Questions

### What batch size does the Macro calendar outbox use for processing rows?

The drain loop processes a maximum of **50 rows** per table per iteration. This `BATCH_SIZE` constant prevents the outbox from overwhelming downstream SQS consumers while maintaining high throughput.

### How does the outbox prevent duplicate calendar invite messages?

The system acquires a **row-level lock** using `FOR UPDATE OF … SKIP LOCKED` before processing each row. After successfully enqueueing the message to SQS, it sets the `published_at` timestamp. If a worker crashes after publishing but before updating the database, the row remains locked until the transaction times out, ensuring another worker cannot process it concurrently.

### What happens when calendar sync is temporarily disabled?

When the `calendar_sync_enabled` flag is false, the drain loop skips the scheduler invocation but continues to process any existing rows already present in the outbox tables. Previously queued jobs are not lost; they remain in the table with `NULL` `published_at` values until sync is re-enabled or an operator manually triggers republishing.

### How are malformed or unknown calendar job types handled?

If the outbox encounters a row with an unrecognized `kind` value during message mapping, it immediately marks the row as published by setting `published_at` to the current timestamp. This prevents the malformed entry from blocking subsequent valid jobs while allowing operators to audit the failure through logging and database inspection.