Architecture of the Campaign Manager and Message Queue in Listmonk

Listmonk's mailing engine orchestrates bulk delivery through a Manager-Pipe-Worker architecture that decouples campaign scheduling from message transmission via buffered Go channels.

The knadh/listmonk repository implements a high-throughput campaign management system designed to handle millions of subscribers without overwhelming upstream providers. The architecture of the campaign manager and message queue separates campaign orchestration from message delivery through an internal queue system that supports configurable concurrency, rate limiting, and multiple messenger back-ends.

Core Components

Three primary structures coordinate the campaign execution pipeline: the Manager maintains global state and worker pools, Pipe instances handle per-campaign logic, and Workers consume messages from buffered channels.

The Manager (internal/manager/manager.go)

The Manager struct serves as the central orchestrator, initialized via manager.New() with a configuration object, database store, i18n instance, and logger. It holds two primary buffered channels: campMsgQ for campaign-bound messages and msgQ for ad-hoc messages. The Manager maintains a registry of messenger implementations and a cache of compiled templates.

When cfg.ScanCampaigns is enabled, the Manager starts a background ticker that periodically invokes scanCampaigns(). This method queries the database via Store.NextCampaigns() and creates a new pipe for each active campaign, pushing it onto the nextPipes channel for execution.

The Pipe (internal/manager/pipe.go)

A pipe represents a single running campaign instance, created internally via newPipe(). Each pipe maintains:

  • A sync.WaitGroup tracking in-flight messages
  • An atomic counter for sent messages and errors
  • A rate limiter enforcing per-minute quotas via cfg.MessageRate
  • A stop flag for pausing or canceling campaigns

The pipe fetches subscriber batches through Store.NextSubscribers() and generates CampaignMessage objects by calling Manager.NewCampaignMessage() to render templates with tracking links. Messages are enqueued onto Manager.campMsgQ with sliding-window throttling to respect provider rate limits.

When the WaitGroup reaches zero (all messages processed), the pipe executes cleanup(), which updates campaign statistics in the database, changes the campaign status to finished or paused, and triggers admin notifications via Manager.fnNotify.

The Workers

The Manager spawns cfg.Concurrency goroutines running the worker() method. These workers continuously select from campMsgQ and msgQ, converting CampaignMessage objects into models.Message structs by adding tracking headers, unsubscribe links, and attachments.

Workers invoke Messenger.Push() on the appropriate back-end (SMTP, SES, Postback, etc.) and update pipe counters via wg.Done(), rate.Incr(), and sent.Add(). They enforce cfg.MessageRate limits using a sliding window algorithm to prevent upstream throttling.

Data Flow Through the Pipeline

Campaign execution follows a four-stage pipeline that decouples database polling from network I/O:

  1. Scan Loop – The Manager's scanCampaigns ticker (configured via cfg.ScanInterval) polls Store.NextCampaigns() and initializes a dedicated pipe for each ready campaign.

  2. Pipe Execution – For each pipe, NextSubscribers() pulls batches of size cfg.BatchSize, renders personalized content through the template engine, and pushes CampaignMessage objects onto the buffered campMsgQ channel.

  3. Worker Consumption – Concurrent workers dequeue messages, construct final models.Message instances (including MIME headers and tracking pixels), and deliver them via the registered messenger interface located in internal/messenger/.

  4. Completion – When a pipe exhausts its subscriber list and all workers call wg.Done(), the cleanup() routine persists final metrics, updates campaign status, and sends completion alerts through the notifs package.

Rate Limiting and Error Handling

The architecture implements multiple safeguards to ensure deliverability and prevent abuse.

Message Rate Control: Workers respect cfg.MessageRate (messages per second) and optional sliding-window throttling configured via cfg.SlidingWindow parameters. The rate counter implemented in internal/utils/utils.go provides atomic increment operations with time-windowed reset logic.

Error Thresholds: The pipe.OnError method increments an atomic error counter for each delivery failure. Once errors exceed cfg.MaxSendErrors, the pipe automatically stops and marks the campaign as paused, preventing reputation damage from repeated bounces.

Extending the Architecture

New message transports integrate via the Messenger interface requiring four methods: Name(), Push(), Flush(), and Close(). Implementations register through Manager.AddMessenger(), making them available for worker dispatch.

Template customization occurs through TemplateFuncs() and custom helper functions that extend the Sprig template library, allowing custom formatting logic in campaign templates.

Implementation Example

The following snippet demonstrates initializing the manager, registering an SMTP messenger, and pushing ad-hoc messages:

// Initialize the manager (normally done in cmd/main.go)
cfg := manager.Config{
    BatchSize:       500,
    Concurrency:     4,
    MessageRate:     200,
    ScanCampaigns:   true,
    ScanInterval:    30 * time.Second,
    // ...other fields...
}
mgr := manager.New(cfg, store, i18nInstance, log.New(os.Stdout, "", log.LstdFlags))

// Register an e‑mail messenger (implementation in internal/messenger/email)
emailMessenger := email.NewSMTP(...)
// Add messenger to manager
if err := mgr.AddMessenger(emailMessenger); err != nil {
    log.Fatalf("messenger error: %v", err)
}

// Start the manager (runs scan loop and workers)
go mgr.Run()

// Push a one‑off message (outside of a campaign)
msg := models.Message{
    From:    "noreply@example.com",
    To:      []string{"user@example.com"},
    Subject: "Welcome!",
    Body:    []byte("Hello there!"),
}
if err := mgr.PushMessage(msg); err != nil {
    log.Printf("push error: %v", err)
}

// Graceful shutdown (e.g., on SIGTERM)
mgr.Close()

Key Source Files

Understanding the full implementation requires examining these specific files in the knadh/listmonk repository:

Summary

  • Manager-Pipe-Worker architecture decouples campaign scheduling from message delivery using Go channels
  • Buffered channels (campMsgQ, msgQ) allow workers to consume messages asynchronously from multiple concurrent pipes
  • Per-campaign pipes maintain isolated state including rate limiting, error counting, and WaitGroup coordination for graceful shutdown
  • Configurable concurrency via cfg.Concurrency and cfg.MessageRate controls throughput to prevent upstream throttling
  • Pluggable messengers implement the Messenger interface to support Email, SMS, Webhooks, and custom transports

Frequently Asked Questions

How does Listmonk handle concurrent campaign execution?

The Manager spawns a fixed pool of worker goroutines (configured via cfg.Concurrency) that consume from shared channels. Each active campaign runs in its own pipe instance that enqueues messages onto campMsgQ. Workers pull from this channel indiscriminately, allowing multiple campaigns to interleave their messages while maintaining separate rate limits and error counters per pipe.

What happens when a campaign exhausts its subscriber list?

When a pipe finishes fetching all subscribers via Store.NextSubscribers(), it waits for all in-flight messages to complete using wg.Wait(). The cleanup() method then updates the database with final sent/error counts, changes the campaign status to finished or paused, and dispatches a notification to administrators through the notifs package.

How does the message queue prevent overwhelming SMTP providers?

Workers enforce cfg.MessageRate (messages per second) using an atomic rate counter with sliding-window support. Additionally, pipes implement sliding-window throttling before enqueueing messages, ensuring the system respects both per-campaign and global rate limits before messages even reach the worker pool.

Can the campaign manager handle ad-hoc messages outside of campaigns?

Yes. The Manager exposes PushMessage() which sends models.Message objects directly to the msgQ channel, bypassing the campaign pipe system. This allows transactional emails or system notifications to share the same worker pool and messenger infrastructure while maintaining separate priority handling from bulk campaign traffic.

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 →