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

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 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

The core orchestration logic resides in 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, 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.

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. 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:

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:

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 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.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →