# Architecture of the Campaign Manager and Message Queue in Listmonk

> Discover Listmonk's Manager-Pipe-Worker architecture. Learn how it decouples campaign scheduling from message transmission using buffered Go channels for efficient bulk delivery.

- Repository: [Kailash Nadh/listmonk](https://github.com/knadh/listmonk)
- Tags: architecture
- Published: 2026-05-19

---

**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`](https://github.com/knadh/listmonk/blob/main/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:

```go
// 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:

- **[`internal/manager/manager.go`](https://github.com/knadh/listmonk/blob/main/internal/manager/manager.go)** – Core manager struct, configuration, channel setup, scanning loop, and worker implementation
- **[`internal/manager/pipe.go`](https://github.com/knadh/listmonk/blob/main/internal/manager/pipe.go)** – Per-campaign pipe handling including subscriber fetching, message creation, error handling, and cleanup
- **[`models/messages.go`](https://github.com/knadh/listmonk/blob/main/models/messages.go)** – Definition of the `Message` struct that workers ultimately push to messengers
- **[`internal/messenger/email/email.go`](https://github.com/knadh/listmonk/blob/main/internal/messenger/email/email.go)** – Reference implementation of the `Messenger` interface for SMTP delivery
- **[`internal/notifs/notifs.go`](https://github.com/knadh/listmonk/blob/main/internal/notifs/notifs.go)** – Notification helper for campaign status alerts
- **[`internal/core/campaigns.go`](https://github.com/knadh/listmonk/blob/main/internal/core/campaigns.go)** – Campaign model methods and status constants
- **[`internal/utils/utils.go`](https://github.com/knadh/listmonk/blob/main/internal/utils/utils.go)** – Rate-counter wrapper and other utility helpers

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