Client SDK

A single Go package for everything DSend does: publish, consume, acknowledge, manage queues and exchanges, and read broker metrics.

Overview

The SDK lives in github.com/Ali-Hasan-Khan/dsend/client. It exposes two clients that mirror the broker's two roles:

  • Producer — publishes messages, manages queues, and manages exchanges and bindings.
  • Consumer — subscribes to a queue, receives deliveries, and acknowledges them.

Publishing routes a message through an exchange to every bound queue whose binding key matches the publish's routing key; consuming subscribes to a queue by name. See Exchanges & routing.

Each client holds a single persistent TCP connection to the broker and speaks the JSON wire protocol directly. No HTTP, no auto-reconnect, no magic — just synchronous calls and explicit context.Context handling.

Addresses. The broker listens on 127.0.0.1:8080 by default. Pass that as the addr argument to NewProducer and NewConsumer. The CLI hard-codes the same address.

Quickstart

1. Start the broker

Build and run the server. It binds 127.0.0.1:8080 and appends to ./data/wal.log.

terminal
$ go build -o dsend ./cmd/dsend
$ ./dsend server

2. Publish

Create the queue, bind it to an exchange, then publish with a routing key. See Exchanges & routing.

main.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)
	}

	if err := p.Publish(context.Background(), "default", "orders", "order-1042 created", 0); err != nil {
		panic(err)
	}
	fmt.Println("published")
}

3. Consume and acknowledge

main.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)
		}

		// do real work here...
		fmt.Println("received:", msg.Payload)

		if err := c.Ack(msg.AckToken); err != nil {
			panic(err)
		}
	}
}
Don't skip the ack. DSend guarantees at-least-once delivery. Until you call Ack(msg.AckToken), the message stays in the broker's in-flight registry and will be redelivered if you take too long. See Delivery semantics.

Producer

A producer dials the broker once and keeps the connection open. All methods take a context.Context so you can time out or cancel cleanly.

Construct

NewProducer(addr string) (*Producer, error)

Opens a persistent TCP connection to the broker at addr.

Methods

MethodDescription
Publish(ctx, exchangeName, routingKey, payload, ttl) Routes a message through exchangeName to every bound queue whose binding matches routingKey. ttl is the message's time-to-live as a time.Duration; pass 0 for no expiry. Returns no route found when nothing matches.
CreateQueue(ctx, name) Creates a named queue.
DeleteQueue(ctx, name) Deletes a named queue and its bindings.
ListQueues(ctx) ([]string, error) Returns the names of all queues on the broker.
CreateExchange(ctx, name, exchangeType) Creates a named exchange. exchangeType is direct, fanout, or topic.
DeleteExchange(ctx, name) Deletes an exchange. Fails while it still has bindings.
ListExchanges(ctx) ([]string, error) Returns the names of all exchanges on the broker.
BindQueue(ctx, exchangeName, queueName, bindingKey) Routes messages whose routing key matches bindingKey to queueName.
UnbindQueue(ctx, exchangeName, queueName, bindingKey) Removes a binding. Publishing with that key then returns no route found.
Metrics(ctx) (*BrokerMetrics, error) Returns broker-wide metrics plus a breakdown per queue.
QueueMetrics(ctx, queueName) (*QueueMetric, error) Returns metrics for a single queue.
Close() error Closes the underlying connection. Always defer this.

Publishing

Payloads are strings. Keep the payload small and treat it as opaque — the broker never inspects the contents. A message is written to the WAL before it is queued, so a crash between the two steps doesn't lose it.

Bind queues before publishing. Only the default exchange is pre-bound to the default queue at broker startup. A publish whose routing key matches no binding returns no route found. Create the queue and BindQueue it to an exchange before publishing to it.
publish.go
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()

// Named queue — create it, bind it, then publish.
if err := p.CreateQueue(ctx, "emails"); err != nil {
	panic(err)
}
if err := p.BindQueue(ctx, "default", "emails", "emails"); err != nil {
	panic(err)
}
err := p.Publish(ctx, "default", "emails", `{"to":"ada@example.com","subject":"hi"}`, 0)

// Default queue — pre-bound to the default exchange at startup.
err = p.Publish(ctx, "default", "default", "hello, world", 0)
Blocking on full queues. Each queue is a fixed-size ring buffer (capacity 20 by default). When the queue is full, Publish blocks until a consumer frees space or the context is cancelled.

Queue administration

admin.go
if err := p.CreateQueue(ctx, "orders"); err != nil {
	panic(err)
}

names, err := p.ListQueues(ctx)
if err != nil {
	panic(err)
}
fmt.Println(names) // ["default", "orders"]

if err := p.DeleteQueue(ctx, "orders"); err != nil {
	panic(err)
}

Consumer

A consumer is push-based: subscribe once, then loop on Receive. The broker schedules deliveries round-robin across all consumers subscribed to a queue, up to each consumer's prefetch limit (default 10 unacked deliveries). See Delivery semantics.

Construct & subscribe

NewConsumer(addr string) (*Consumer, error)
Subscribe(queueName string) error

Subscribe registers this connection as a consumer session for the queue. One consumer connection can subscribe to a single queue at a time.

Receiving

Receive(ctx context.Context) (*ReceivedMessage, error)

Blocks until the broker pushes the next delivery. A ReceivedMessage carries:

FieldTypeMeaning
IDstringUnique message identifier.
PayloadstringOpaque message body.
AckTokenstringProof of delivery; pass it to Ack.
Receive timeouts. Receive sets a 10-second read deadline on the socket and returns an error if no delivery arrives in time. In a real consumer, retry Receive in a loop instead of exiting.

Acknowledging

Ack(token string) error

Tell the broker the message was processed successfully. The broker removes the ack token from its in-flight registry and never redelivers that message.

consumer.go
for {
	msg, err := c.Receive(ctx)
	if err != nil {
		// transient read deadline or network error; try again
		if ctx.Err() != nil {
			return // caller cancelled
		}
		continue
	}

	if err := handle(msg.Payload); err != nil {
		// skip ack: the broker will redeliver later
		continue
	}

	if err := c.Ack(msg.AckToken); err != nil {
		return
	}
}

Unsubscribing

Unsubscribe() error

Removes this session from the queue's round-robin rotation and stops deliveries. Closing the connection also unregisters the session.

Managing queues

A queue is an isolated, fixed-capacity in-memory ring buffer with its own consumers, in-flight messages, DLQ, and metrics. The default queue is created at broker startup and pre-bound to the default exchange; every other queue must be created and bound to an exchange before it can receive messages.

  • Auto-created — the default queue exists as soon as the broker starts, bound to the default exchange.
  • Explicit creation — CreateQueue or dsend queue create for named queues.
  • Binding — a queue only receives messages routed to it by an exchange binding. See Exchanges & routing.
  • Capacity — 20 messages per queue by default (see configuration).
  • Deletion — DeleteQueue removes the queue and its bindings from the broker.
cli
$ ./dsend queue create orders
Queue created successfully

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

$ ./dsend queue list
Queues: [default orders]

$ ./dsend queue delete orders
Queue deleted successfully

Message expiry

A message can be given a time-to-live (TTL) at publish time. The broker anchors expiry to its own clock: when a message is accepted it computes ExpiryAt = now + ttl. Once that time passes, the message is considered expired and is dead-lettered — moved to the queue's DLQ and reflected in DlqCount.

Expiry is enforced at delivery time, so an expired message is never handed to a consumer. A dedicated expiry worker additionally sweeps each queue on a timer (ExpiryInterval, default 1s), so expired messages are also cleared when no consumer is connected or the broker is otherwise idle.

Publishing with a TTL

Pass a time.Duration as the final argument to Publish. 0 (or any non-positive value) means the message never expires.

publish_with_ttl.go
// expires 10 seconds after the broker accepts it
err := p.Publish(ctx, "default", "orders", "order-1042 created", 10*time.Second)

// never expires
err = p.Publish(ctx, "default", "orders", "order-1042 created", 0)
cli
$ ./dsend publish --exchange default --routing-key orders --ttl 10s "order-1042 created"
Message Sent successfully

--ttl accepts Go duration strings (10s, 5m, 1h30m). All flags (--exchange, --routing-key, --ttl) must appear before the <payload> argument — Go flag parsing stops at the first positional. Messages published without --ttl never expire.

Exchanges & routing

A queue never receives a message directly. Producers publish to an exchange, and the exchange routes the message to every queue bound to it whose binding key matches the publish's routing key.

Exchange types

TypeRouting ruleBest for
direct Routing key must exactly match the binding key. One-to-one task routing — a key per queue.
fanout Every bound queue receives every message; the routing key is ignored. Broadcasts and event fan-out.
topic Binding keys support * (one dot-separated segment) and # (zero or more segments). Pattern-based routing across event namespaces.

Default exchange

The broker auto-creates a default exchange of type direct at startup and pre-binds the default queue to it with binding key default. The default exchange is reserved: it can't be created again or deleted.

Binding keys

A binding key is non-empty, at most 255 characters, and dot-separated into segments. Each segment is letters, digits, _, or -. On topic exchanges, * matches a single segment and # matches zero or more but only as the final segment. Invalid keys (empty, empty segments, or a misplaced #) are rejected with binding key is invalid.

exchanges.go
// a topic exchange for event namespaces
if err := p.CreateExchange(ctx, "events", "topic"); err != nil {
	panic(err)
}

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

// "orders.*" routes order events to the orders queue;
// "events.#" catches every event for the audit queue.
if err := p.BindQueue(ctx, "events", "orders", "orders.*"); err != nil {
	panic(err)
}
if err := p.BindQueue(ctx, "events", "audit", "events.#"); err != nil {
	panic(err)
}

if err := p.Publish(ctx, "events", "orders.created", "order-1042 created", 0); err != nil {
	panic(err)
}
if err := p.Publish(ctx, "events", "orders.shipped", "order-1042 shipped", 0); err != nil {
	panic(err)
}
No route is an error. Publishing when no binding matches returns no route found. After UnbindQueue, publishes for that key fail instead of silently disappearing.
Deleting exchanges. DeleteExchange refuses to remove an exchange that still has bindings. Unbind or delete its bound queues first.
cli
$ ./dsend exchange create events topic
Exchange created successfully

$ ./dsend queue create orders
Queue created successfully

$ ./dsend queue bind events orders "orders.*"
Queue bounded successfully

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

$ ./dsend exchange list
Exchanges: [default events]

Metrics

Every counter maps to a concrete delivery event. Query broker-wide metrics or scope them to one queue. The CLI exposes the same data: ./dsend metrics or ./dsend metrics --queue orders.

Types

internal/model/metrics.go
type Metric struct {
	ProducedCount        int
	AckedCount           int
	DlqCount             int
	InflightCount        int
	RedeliveredCount     int
	ConsumerSessionCount int
	QueueDepth           int
}

type QueueMetric struct {
	Name string
	Metric
}

type BrokerMetrics struct {
	Queues []QueueMetric
	Total  Metric
}

Fields

FieldMeaning
ProducedCountMessages published to the queue.
AckedCountMessages acknowledged by consumers.
QueueDepthMessages currently waiting in the ring buffer.
InflightCountMessages delivered but not yet acknowledged.
RedeliveredCountMessages requeued after an ack timeout.
DlqCountMessages dead-lettered after exhausting retries or expiring.
ConsumerSessionCountConnected consumers subscribed to the queue.
metrics.go
m, err := p.Metrics(ctx)
if err != nil {
	panic(err)
}

for _, q := range m.Queues {
	fmt.Printf("%s: depth=%d in-flight=%d acked=%d\n",
		q.Name, q.QueueDepth, q.InflightCount, q.AckedCount)
}

q, err := p.QueueMetrics(ctx, "orders")
if err != nil {
	panic(err)
}
fmt.Println(q.DlqCount)

CLI reference

The dsend binary wraps the SDK for scripting and quick tests.

CommandDescription
dsend server Starts the broker on 127.0.0.1:8080 with the WAL at ./data/wal.log.
dsend publish [--exchange <name>] [--routing-key <key>] [--ttl <duration>] "<message>" Publishes a message through the named exchange under the given routing key (default default). Defaults to the default exchange. --ttl sets the message's time-to-live (e.g. 10s, 5m). All flags must precede the positional <message> argument.
dsend subscribe [--queue <name>] Subscribes, receives messages, and acks each one. Ctrl-C to stop.
dsend exchange create <name> <type> Creates an exchange (direct, fanout, or topic).
dsend exchange delete <name> Deletes an exchange that has no bindings.
dsend exchange list Lists all exchange names.
dsend metrics [--queue <name>] Prints broker metrics. Scopes to one queue with --queue.
dsend queue create <name> Creates a named queue.
dsend queue delete <name> Deletes a named queue.
dsend queue list Lists all queue names.
dsend queue bind <exchange> <queue> <bindingKey> Routes matching messages from the exchange to the queue.
dsend queue unbind <exchange> <queue> <bindingKey> Removes a binding.
The CLI targets localhost:8080 and the broker must be running first. There's no persistence across dsend server restarts beyond what the WAL recovers.

Wire protocol

Clients and broker exchange length-delimited JSON over TCP. Each message is a request with a type field; the broker replies with a response carrying success and, on failure, error.

Request types

TypeSent byPurpose
publishProducerRoute a message through an exchange to bound queues.
subscribeConsumerRegister a consumer session for a queue.
ackConsumerAcknowledge a delivery by token.
unsubscribeConsumerLeave the queue's round-robin rotation.
metricsProducerFetch broker-wide or per-queue metrics.
create_queueProducerCreate a named queue.
delete_queueProducerDelete a named queue.
list_queuesProducerList all queue names.
bind_queueProducerBind a queue to an exchange with a binding key.
unbind_queueProducerRemove a binding.
create_exchangeProducerCreate an exchange with a type.
delete_exchangeProducerDelete an exchange without bindings.
list_exchangesProducerList all exchange names.

Example

client → broker
{"type":"publish","exchange":"default","routing_key":"orders","message":{"payload":"order-1042 created"}}
client → broker (publish with TTL)
{"type":"publish","exchange":"default","routing_key":"orders","message":{"payload":"order-1042 created"},"ttl":10000000000}

ttl is a time.Duration and is serialized as an integer number of nanoseconds (10s = 10000000000). The broker converts it to an absolute expiry_at on its own clock. Omit it for no expiry.

client → broker (bind)
{"type":"bind_queue","exchange":"default","queue":"orders","binding_key":"orders"}
broker → client
{"success":true,"error":""}
broker → consumer (delivery)
{
  "success": true,
  "message": {
    "id": "a1b2c3d4-...",
    "payload": "order-1042 created",
    "timestamp": "2026-07-31T12:00:00Z",
    "expiry_at": "2026-07-31T12:00:10Z",
    "retry": 0
  },
  "ack_token": "9f8e7d6c-..."
}

Delivery semantics

Push with round-robin & prefetch

The broker pushes deliveries to subscribed consumers in round-robin order. Add more consumers to a queue and load is spread across them automatically. Each consumer can hold up to ConsumerPrefetch unacknowledged deliveries at once (default 10).

A consumer at its prefetch limit is skipped by the scheduler until it acknowledges a delivery. This is backpressure: a slow consumer isn't flooded with more work than it can handle, and a fast consumer isn't starved behind a slow one. Messages bound for a saturated consumer stay queued instead of being pushed.

At-least-once

When a message is delivered, the broker moves it from the ring buffer into an in-flight registry keyed by ack token. It stays there until one of three things happens:

  1. Ack — the consumer processed it; it's removed permanently.
  2. Ack timeout — if the consumer doesn't ack within AckTimeout, the message is requeued and redelivered.
  3. Retry limit — after MaxRetries redeliveries, the message is moved to the queue's DLQ.
Ordering. A message can be delivered more than once (at-least-once, not exactly-once). Make your consumers idempotent when duplicate processing matters.

Default configuration

SettingDefaultEffect
AckTimeout100 sHow long a delivery may sit unacked before redelivery.
RedeliveryInterval5 sHow often expired in-flight messages are swept.
ExpiryInterval1 sHow often the expiry worker sweeps queues for TTL-expired messages.
MaxRetries3Redeliveries before a message is dead-lettered.
QueueSize20Capacity of each queue's ring buffer.
ConsumerPrefetch10Max unacknowledged deliveries a single consumer may hold in flight.

Configuration is set in engine.DefaultConfig() at broker startup (cmd/dsend/server.go).

Dead-letter queue

Each queue has its own DLQ. Messages that exhaust their retries — or expire past their TTL — land there instead of being dropped. The DlqCount metric tracks how many have.

Persistence & recovery

The broker appends every state change — publish, ack, requeue, dead-letter, and exchange and binding changes — to a write-ahead log. On restart it replays the WAL and rebuilds queues, exchanges, bindings, in-flight state, and session ownership.