Skip to content

About

Single-node message broker built from scratch with .NET 10. Features TCP-based messaging, exchange routing, persistent queues, acknowledgements, publisher confirms, and filesystem-backed storage.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Repository files navigation

SystemDesign.MessageQueue

Single-node message broker built from scratch with .NET 10.

The broker provides an HTTP management plane and a persistent TCP data plane with a custom binary protocol. It does not depend on an external database or an existing message broker.

Features

  • Persistent TCP connections for producers and consumers
  • Custom binary framed protocol
  • Direct, Fanout and Topic exchanges
  • Queues and bindings
  • Routing by routing key
  • Persistent and transient messages
  • Publisher confirms
  • Manual ACK/NACK
  • Message redelivery
  • Consumer prefetch
  • Message and queue TTL
  • Dead Letter Exchange (DLX)
  • Durable topology
  • Filesystem-backed append-only message storage
  • Checkpoints and recovery
  • Heartbeat with Ping/Pong
  • Bounded channels and backpressure
  • HTTP management API
  • Swagger/OpenAPI
  • .NET Client SDK
  • Producer and Consumer samples

Architecture

flowchart LR
    Producer[Producer<br/>Client SDK]
    Consumer[Consumer<br/>Client SDK]
    Admin[Administrator]

    TCP[Broker TCP Server<br/>:5673]
    API[Management API<br/>HTTP / Swagger]

    Topology[Topology Registry<br/>Exchanges / Queues / Bindings]
    Routing[Routing Engine]
    Queue[Broker Queue]
    Delivery[Delivery Dispatcher]

    Segments[Append-only<br/>Queue Segments]
    Journal[Topology Journal]
    Checkpoints[Checkpoints]

    Producer -->|Publish / TCP| TCP
    TCP -->|Delivery / TCP| Consumer

    Admin -->|HTTP| API
    API --> Topology

    TCP --> Routing
    Topology --> Routing
    Routing --> Queue
    Queue --> Delivery
    Delivery --> TCP

    Queue --> Segments
    Topology --> Journal
    Segments --> Checkpoints
Loading

The solution is split by runtime responsibility:

src/
├── SystemDesign.MessageQueue.Api
├── SystemDesign.MessageQueue.Broker
├── SystemDesign.MessageQueue.Storage
├── SystemDesign.MessageQueue.Protocol
└── SystemDesign.MessageQueue.Client

samples/
├── SystemDesign.MessageQueue.Sample.Producer
└── SystemDesign.MessageQueue.Sample.Consumer

Projects

SystemDesign.MessageQueue.Api

Application host and composition root. Provides the HTTP management API, Swagger, TCP listener, and background services.

SystemDesign.MessageQueue.Broker

Core broker logic: topology, exchanges, queues, bindings, routing, consumers, delivery, ACK/NACK, prefetch, publisher confirms, TTL, and dead lettering.

SystemDesign.MessageQueue.Storage

Filesystem persistence: topology journal, append-only queue segments, checksums, checkpoints, recovery, and cleanup.

SystemDesign.MessageQueue.Protocol

Binary protocol framing, commands, and wire contracts shared by the broker and client.

SystemDesign.MessageQueue.Client

.NET client SDK that hides TcpClient, NetworkStream, framing, correlation IDs, subscriptions, ACK/NACK, and connection management.

Message Flow

sequenceDiagram
    participant P as Producer
    participant B as Broker
    participant S as Storage
    participant Q as Queue
    participant C as Consumer

    P->>B: Publish(exchange, routingKey)
    B->>B: Route message
    B->>S: Append persistent message
    S-->>B: Stored
    B->>Q: Enqueue
    B-->>P: PublishConfirm

    Q->>C: Delivery
    C-->>B: ACK
    B->>S: Checkpoint completion
Loading

A producer publishes a message to an exchange using a routing key. The routing engine resolves matching bindings and determines the destination queues.

Persistent messages are written to append-only storage before the corresponding publisher confirmation according to the configured durability mode.

Consumers receive messages over persistent TCP connections. With manual acknowledgement enabled, delivered messages remain unacknowledged until the consumer sends ACK or NACK.

Exchanges and Routing

The broker supports three exchange types.

Direct

A message is routed when its routing key exactly matches the binding routing key.

orders
   │
   │ order.created
   ▼
billing-queue

Fanout

The message is routed to every queue bound to the exchange. The routing key is ignored.

Topic

Topic exchanges support wildcard routing:

  • * matches exactly one segment
  • # matches zero or more segments

Examples:

order.*
order.#
*.created

TCP Protocol

Producer and consumer traffic does not use HTTP.

Clients maintain persistent TCP connections to the broker. The default TCP port is:

5673

Messages are transferred using a custom binary framed protocol.

Conceptually, each frame contains:

Version | Command | Flags | CorrelationId | PayloadLength | Payload

Supported protocol operations include:

Connect
Publish
PublishConfirm
PublishReject
Subscribe
Delivery
Ack
Nack
Ping
Pong
Error

The frame reader handles TCP fragmentation and coalescing and validates frame sizes before allocating payload buffers.

Each connection uses a bounded outgoing channel and a dedicated send loop, preventing concurrent writes from corrupting the TCP stream.

Persistence and Recovery

Persistent messages are stored on the local filesystem without an external database.

Each durable queue uses rotating append-only segment files:

data/
├── topology/
│   └── topology.log
│
└── queues/
    ├── orders/
    │   ├── 00000000000000000000.log
    │   ├── 00000000000000100000.log
    │   └── checkpoint.dat
    │
    └── notifications/
        ├── 00000000000000000000.log
        └── checkpoint.dat

Records contain message metadata, payload information, and CRC32 checksums.

ACK does not rewrite records already stored inside a segment. Completed messages are tracked separately using checkpoints.

At startup, the broker:

  1. Replays the durable topology journal.
  2. Discovers queue segment files.
  3. Scans and validates stored records.
  4. Stops at the last valid boundary if a partial or corrupt trailing record is found.
  5. Loads queue checkpoints.
  6. Restores persistent messages that were not completed.
  7. Makes recovered messages available for redelivery.

Durable topology and persistent messages can therefore survive broker restarts.

ACK, NACK and Redelivery

With manual acknowledgement, a message moves through the following states:

Ready
  │
  ▼
Delivered
  │
  ▼
Unacked
  │
  ├── ACK ─────────────► Completed
  │
  └── NACK
       │
       ├── requeue=true ──► Ready
       │
       └── requeue=false ─► DLX / Discard

Prefetch limits the number of outstanding unacknowledged messages for each consumer.

Messages can be returned to the queue for redelivery when processing fails or a delivery is negatively acknowledged with requeue enabled.

Publisher Confirms

Each publish operation contains a correlation ID.

Multiple publish operations can be in flight simultaneously. The broker responds with either:

PublishConfirm

or:

PublishReject

using the corresponding correlation ID.

The client SDK correlates these responses with the original asynchronous publish operations.

TTL and Dead Lettering

Queues can define:

  • message TTL
  • maximum queue length
  • dead-letter exchange
  • dead-letter routing key

Expired or rejected messages can be routed through the configured Dead Letter Exchange using the same routing engine as ordinary messages.

Management API

HTTP is used for broker management and inspection, not for message publishing or delivery.

Available operations include:

GET  /api/topology

GET  /api/exchanges
PUT  /api/exchanges/{name}

GET  /api/queues
GET  /api/queues/{name}
PUT  /api/queues/{name}

POST /api/bindings

Queue information includes runtime counters such as:

ReadyCount
UnackedCount
ConsumerCount

Swagger UI is available at:

/swagger

Client SDK

Producer

await using var connection =
    await MessageQueueConnection.ConnectAsync("localhost", 5673);

var producer = connection.CreateProducer();

var result = await producer.PublishJsonAsync(
    exchange: "orders",
    routingKey: "order.created",
    value: order,
    persistent: true);

Console.WriteLine($"Confirmed: {result.Confirmed}");

Consumer

await using var connection =
    await MessageQueueConnection.ConnectAsync("localhost", 5673);

var consumer = connection.CreateConsumer();

await consumer.SubscribeAsync<Order>(
    queue: "billing",
    prefetch: 10,
    async context =>
    {
        await ProcessAsync(context.Message);
        await context.AckAsync();
    });

await Task.Delay(Timeout.Infinite);

Running

Requirements:

  • .NET 10 SDK

Start the broker:

dotnet run --project src/SystemDesign.MessageQueue.Api

Run the sample consumer:

dotnet run --project samples/SystemDesign.MessageQueue.Sample.Consumer

Run the sample producer:

dotnet run --project samples/SystemDesign.MessageQueue.Sample.Producer

The producer and consumer communicate with the broker through independent persistent TCP connections.

Scope

The project intentionally focuses on the internals of a single-node message broker.

It does not implement:

  • clustering
  • replication
  • consensus
  • distributed queue ownership
  • cross-node routing
  • leader election
  • service discovery
  • AMQP
  • Kubernetes integration
  • external state storage

Consistency, persistence, and recovery semantics are explicitly designed for a single broker node.

The project focuses on message broker internals: routing, delivery, acknowledgements, persistent TCP connections, binary framing, durable storage, recovery, backpressure, and client/server coordination.


SystemDesign.MessageQueue — Русский

Одноузловой брокер сообщений, реализованный с нуля на .NET 10.

Брокер предоставляет HTTP API для управления и постоянный TCP-контур передачи сообщений с собственным бинарным протоколом. Для работы не требуется внешняя база данных или готовый message broker.

Возможности

  • Постоянные TCP-соединения producer и consumer
  • Собственный бинарный framed-протокол
  • Exchanges типов Direct, Fanout и Topic
  • Очереди и bindings
  • Маршрутизация по routing key
  • Persistent и transient сообщения
  • Publisher confirms
  • Ручные ACK/NACK
  • Повторная доставка сообщений
  • Consumer prefetch
  • TTL сообщений и очередей
  • Dead Letter Exchange (DLX)
  • Durable-топология
  • Файловое append-only хранилище сообщений
  • Checkpoints и восстановление состояния
  • Heartbeat через Ping/Pong
  • Bounded channels и backpressure
  • HTTP Management API
  • Swagger/OpenAPI
  • .NET Client SDK
  • Примеры Producer и Consumer

Архитектура

flowchart LR
    Producer[Producer<br/>Client SDK]
    Consumer[Consumer<br/>Client SDK]
    Admin[Administrator]

    TCP[Broker TCP Server<br/>:5673]
    API[Management API<br/>HTTP / Swagger]

    Topology[Topology Registry<br/>Exchanges / Queues / Bindings]
    Routing[Routing Engine]
    Queue[Broker Queue]
    Delivery[Delivery Dispatcher]

    Segments[Append-only<br/>Queue Segments]
    Journal[Topology Journal]
    Checkpoints[Checkpoints]

    Producer -->|Publish / TCP| TCP
    TCP -->|Delivery / TCP| Consumer

    Admin -->|HTTP| API
    API --> Topology

    TCP --> Routing
    Topology --> Routing
    Routing --> Queue
    Queue --> Delivery
    Delivery --> TCP

    Queue --> Segments
    Topology --> Journal
    Segments --> Checkpoints
Loading

Решение разделено на проекты по их ответственности:

src/
├── SystemDesign.MessageQueue.Api
├── SystemDesign.MessageQueue.Broker
├── SystemDesign.MessageQueue.Storage
├── SystemDesign.MessageQueue.Protocol
└── SystemDesign.MessageQueue.Client

samples/
├── SystemDesign.MessageQueue.Sample.Producer
└── SystemDesign.MessageQueue.Sample.Consumer

Проекты

SystemDesign.MessageQueue.Api

Хост приложения и composition root. Содержит HTTP Management API, Swagger, TCP listener и фоновые сервисы.

SystemDesign.MessageQueue.Broker

Основная логика брокера: топология, exchanges, queues, bindings, маршрутизация, consumers, доставка, ACK/NACK, prefetch, publisher confirms, TTL и dead lettering.

SystemDesign.MessageQueue.Storage

Файловое хранилище: журнал топологии, append-only сегменты очередей, checksums, checkpoints, восстановление и очистка.

SystemDesign.MessageQueue.Protocol

Описание бинарного протокола, frames, команд и wire-контрактов, общих для брокера и клиента.

SystemDesign.MessageQueue.Client

.NET Client SDK, скрывающий работу с TcpClient, NetworkStream, frames, correlation ID, подписками, ACK/NACK и управлением соединением.

Поток сообщения

sequenceDiagram
    participant P as Producer
    participant B as Broker
    participant S as Storage
    participant Q as Queue
    participant C as Consumer

    P->>B: Publish(exchange, routingKey)
    B->>B: Routing
    B->>S: Append persistent message
    S-->>B: Stored
    B->>Q: Enqueue
    B-->>P: PublishConfirm

    Q->>C: Delivery
    C-->>B: ACK
    B->>S: Checkpoint completion
Loading

Producer публикует сообщение в exchange с определённым routing key. Routing Engine находит подходящие bindings и определяет очереди назначения.

Persistent-сообщения записываются в append-only хранилище до отправки соответствующего publisher confirm в соответствии с выбранным режимом durability.

Consumer получает сообщения через постоянное TCP-соединение. При ручном подтверждении доставленное сообщение остаётся в состоянии unacked до получения ACK или NACK.

Exchanges и маршрутизация

Поддерживаются три типа exchange.

Direct

Routing key сообщения должен точно совпасть с routing key binding.

orders
   │
   │ order.created
   ▼
billing-queue

Fanout

Сообщение направляется во все очереди, связанные с exchange. Routing key игнорируется.

Topic

Topic exchange поддерживает wildcard-маршрутизацию:

  • * соответствует ровно одному сегменту
  • # соответствует нулю или нескольким сегментам

Например:

order.*
order.#
*.created

TCP-протокол

Публикация и доставка сообщений не используют HTTP.

Клиенты поддерживают постоянные TCP-соединения с брокером. TCP-порт по умолчанию:

5673

Для передачи данных используется собственный бинарный framed-протокол.

Концептуально frame имеет следующую структуру:

Version | Command | Flags | CorrelationId | PayloadLength | Payload

Поддерживаются команды:

Connect
Publish
PublishConfirm
PublishReject
Subscribe
Delivery
Ack
Nack
Ping
Pong
Error

FrameReader корректно обрабатывает фрагментацию и объединение TCP-пакетов и проверяет размер frame до выделения памяти под payload.

Каждое соединение использует bounded outgoing channel и отдельный send loop, поэтому конкурентная запись нескольких сообщений не повреждает TCP stream.

Хранение и восстановление

Persistent-сообщения сохраняются непосредственно в файловой системе без внешней базы данных.

Для каждой durable-очереди используются ротируемые append-only сегменты:

data/
├── topology/
│   └── topology.log
│
└── queues/
    ├── orders/
    │   ├── 00000000000000000000.log
    │   ├── 00000000000000100000.log
    │   └── checkpoint.dat
    │
    └── notifications/
        ├── 00000000000000000000.log
        └── checkpoint.dat

Записи содержат метаданные сообщения, информацию о payload и CRC32 checksum.

ACK не изменяет уже записанные данные внутри segment. Завершённые сообщения отслеживаются отдельно с помощью checkpoints.

При запуске брокер:

  1. Восстанавливает durable-топологию из журнала.
  2. Находит segment-файлы очередей.
  3. Сканирует и проверяет сохранённые записи.
  4. При обнаружении повреждённой или частично записанной последней записи останавливается на последней корректной границе.
  5. Загружает checkpoints.
  6. Восстанавливает незавершённые persistent-сообщения.
  7. Возвращает восстановленные сообщения в состояние, доступное для повторной доставки.

Таким образом, durable-топология и persistent-сообщения могут переживать перезапуск брокера.

ACK, NACK и повторная доставка

При ручном подтверждении сообщение проходит следующие состояния:

Ready
  │
  ▼
Delivered
  │
  ▼
Unacked
  │
  ├── ACK ─────────────► Completed
  │
  └── NACK
       │
       ├── requeue=true ──► Ready
       │
       └── requeue=false ─► DLX / Discard

Prefetch ограничивает количество неподтверждённых сообщений, одновременно выданных одному consumer.

При ошибке обработки сообщение может быть возвращено в очередь для повторной доставки.

Publisher Confirms

Каждая операция публикации содержит correlation ID.

Несколько операций Publish могут одновременно находиться в обработке. Брокер отвечает:

PublishConfirm

или:

PublishReject

с соответствующим correlation ID.

Client SDK сопоставляет ответ брокера с исходной асинхронной операцией публикации.

TTL и Dead Letter Exchange

Для очереди могут быть настроены:

  • TTL сообщения
  • максимальный размер очереди
  • Dead Letter Exchange
  • Dead Letter Routing Key

Истёкшие или отклонённые сообщения могут быть направлены через настроенный Dead Letter Exchange с использованием обычного Routing Engine.

Management API

HTTP используется только для управления брокером и просмотра его состояния, а не для публикации или доставки сообщений.

Основные endpoints:

GET  /api/topology

GET  /api/exchanges
PUT  /api/exchanges/{name}

GET  /api/queues
GET  /api/queues/{name}
PUT  /api/queues/{name}

POST /api/bindings

Информация об очередях содержит runtime-счётчики:

ReadyCount
UnackedCount
ConsumerCount

Swagger UI доступен по адресу:

/swagger

Client SDK

Producer

await using var connection =
    await MessageQueueConnection.ConnectAsync("localhost", 5673);

var producer = connection.CreateProducer();

var result = await producer.PublishJsonAsync(
    exchange: "orders",
    routingKey: "order.created",
    value: order,
    persistent: true);

Console.WriteLine($"Confirmed: {result.Confirmed}");

Consumer

await using var connection =
    await MessageQueueConnection.ConnectAsync("localhost", 5673);

var consumer = connection.CreateConsumer();

await consumer.SubscribeAsync<Order>(
    queue: "billing",
    prefetch: 10,
    async context =>
    {
        await ProcessAsync(context.Message);
        await context.AckAsync();
    });

await Task.Delay(Timeout.Infinite);

Запуск

Требования:

  • .NET 10 SDK

Запуск брокера:

dotnet run --project src/SystemDesign.MessageQueue.Api

Запуск Consumer:

dotnet run --project samples/SystemDesign.MessageQueue.Sample.Consumer

Запуск Producer:

dotnet run --project samples/SystemDesign.MessageQueue.Sample.Producer

Producer и Consumer работают как независимые приложения и подключаются к брокеру через отдельные постоянные TCP-соединения.

Ограничения

Проект намеренно сосредоточен на внутреннем устройстве single-node message broker.

Не реализованы:

  • кластеризация
  • репликация
  • consensus
  • распределённое владение очередями
  • маршрутизация между узлами
  • leader election
  • service discovery
  • AMQP
  • интеграция с Kubernetes
  • внешнее хранилище состояния

Механизмы consistency, persistence и recovery рассчитаны на работу одного узла брокера.

Основное внимание уделено внутренним механизмам message broker: маршрутизации, доставке, ACK/NACK, постоянным TCP-соединениям, бинарному framing, persistent storage, восстановлению после перезапуска, backpressure и взаимодействию клиента с брокером.

About

Single-node message broker built from scratch with .NET 10. Features TCP-based messaging, exchange routing, persistent queues, acknowledgements, publisher confirms, and filesystem-backed storage.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages