Queue-based message broker for Go

Reliable messaging, written from scratch.

DSend is a lightweight message broker built in Go — no Kafka, no RabbitMQ, no dependencies. Publish, consume, and acknowledge messages with durable storage, retries, and dead-lettering, all over a plain TCP/JSON protocol.

MIT license · Go 1.25 · github.com/Ali-Hasan-Khan/dsend

dsend — a working broker in one binary
# terminal 1 — start the broker
$ go build -o dsend ./cmd/dsend
$ ./dsend server
[server] listening on 127.0.0.1:8080
[wal]    appending to ./data/wal.log

# terminal 2 — create a queue, bind it to an exchange, publish
# terminal 3 — subscribe and receive it
$ ./dsend queue create orders
Queue created successfully

$ ./dsend queue bind default orders orders
Queue bounded successfully

$ ./dsend publish --exchange default --routing-key orders "order-1042 created"
Message Sent successfully

$ ./dsend subscribe --queue orders
Received message: order-1042 created
What you get

Everything a broker needs. Nothing it doesn't.

Nine core building blocks, implemented from scratch and exposed through a small, typed Go SDK.

Named queues

Isolated in-memory ring buffers. Create named queues with the CLI or SDK — the default queue is pre-bound to the default exchange.

Exchanges & routing

Direct, fanout, and topic exchanges route messages to bound queues by routing key. Bind a queue, then publish to a key.

At-least-once delivery

Delivered messages are tracked in flight until a consumer acknowledges them. Nothing is silently lost.

ACK + redelivery

Messages that aren't acknowledged within the timeout are redelivered, up to a configured retry limit.

Dead-letter queue

Messages that exhaust their retries move to a per-queue DLQ instead of vanishing into the void.

Message expiry (TTL)

Give a message a time-to-live at publish time. Once the TTL elapses, the message is dead-lettered before it can be delivered.

Consumer backpressure

Each consumer holds at most a prefetch count of unacked deliveries. Slow consumers don't hog the queue — fast ones aren't starved.

Durable write-ahead log

Every publish, ack, and requeue is appended to a WAL. Restart the broker and the queues come back.

Runtime metrics

Produced, acked, in-flight, redelivered, DLQ depth, and consumer sessions — per queue or broker-wide.

How it works

Producers push. Consumers pull with a broker in between.

One TCP connection per client, JSON on the wire, and a broker core that schedules delivery round-robin across consumers — up to each consumer's prefetch limit.

Producer

client.NewProducer

  • Publish
  • Metrics
  • Queue admin
  • Exchange admin
TCP · JSON

Broker

per-queue runtime

  • exchanges
  • ring buffer
  • in-flight
  • DLQ
  • WAL
round-robin

Consumer

client.NewConsumer

  • Subscribe
  • Receive
  • Ack

↩ Consumers return an ack token after processing; the broker removes the message from its in-flight registry. Each consumer is capped at a prefetch count of unacked deliveries, and unacked messages are redelivered ack-timeout later.

Quickstart

A producer and a consumer, in two files.

The SDK is a single Go package: github.com/Ali-Hasan-Khan/dsend/client.

producer.go
package main

import (
  "context"
  "fmt"

  "github.com/Ali-Hasan-Khan/dsend/client"
)

func main() {
  p, err := client.NewProducer("localhost:8080")
  if err != nil {
    panic(err)
  }
  defer p.Close()

  if err := p.CreateQueue(context.Background(), "orders"); err != nil {
    panic(err)
  }

  // bind "orders" to the default exchange under routing key "orders"
  if err := p.BindQueue(context.Background(), "default", "orders", "orders"); err != nil {
    panic(err)
  }

  err = p.Publish(context.Background(),
    "default", "orders", "order-1042 created")
  if err != nil {
    panic(err)
  }
  fmt.Println("published")
}
consumer.go
package main

import (
  "context"
  "fmt"

  "github.com/Ali-Hasan-Khan/dsend/client"
)

func main() {
  c, err := client.NewConsumer("localhost:8080")
  if err != nil {
    panic(err)
  }
  defer c.Close()

  if err := c.Subscribe("orders"); err != nil {
    panic(err)
  }

  for {
    msg, err := c.Receive(context.Background())
    if err != nil {
      panic(err)
    }
    fmt.Println("got:", msg.Payload)
    c.Ack(msg.AckToken) // mark as processed
  }
}
1

Start the broker

Build and run the server. It listens on 127.0.0.1:8080 and persists to ./data/wal.log.

2

Run the producer

Open a connection, create the queue (or use default), bind it to the default exchange, then call Publish with a routing key.

3

Consume & ack

Subscribe, receive deliveries, and call Ack after processing. Failures are redelivered automatically.

Observability

See what the broker is doing, at a glance.

Query broker-wide or per-queue metrics through the SDK or the CLI. Every counter tracks a concrete delivery event.

ProducedCount
messages appended to the WAL and queued
AckedCount
messages acknowledged by consumers
QueueDepth / InflightCount
waiting in the ring buffer vs. reserved by a consumer
RedeliveredCount / DlqCount
requeued after ack timeout vs. dead-lettered after retries or expiry
Metrics reference
$ ./dsend metrics --queue orders
Queue: orders

ProducedCount:        10
QueueDepth:            0
InflightCount:         0
DlqCount:              0
ConsumerSessionCount:  1
AckedCount:           10
RedeliveredCount:      0

Ready to move messages between Go services?

Read the SDK docs and get your first producer running in five minutes.

Open the docs