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.
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.
$ 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.
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
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)
}
}
}
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
Opens a persistent TCP connection to the broker at addr.
Methods
| Method | Description |
|---|---|
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.
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.
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)
Publish blocks until a consumer
frees space or the context is cancelled.
Queue administration
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
Subscribe registers this connection as a consumer session for the
queue. One consumer connection can subscribe to a single queue at a time.
Receiving
Blocks until the broker pushes the next delivery. A ReceivedMessage carries:
| Field | Type | Meaning |
|---|---|---|
ID | string | Unique message identifier. |
Payload | string | Opaque message body. |
AckToken | string | Proof of delivery; pass it to Ack. |
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
Tell the broker the message was processed successfully. The broker removes the ack token from its in-flight registry and never redelivers that message.
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
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
defaultqueue exists as soon as the broker starts, bound to thedefaultexchange. - Explicit creation —
CreateQueueordsend queue createfor 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 —
DeleteQueueremoves the queue and its bindings from the broker.
$ ./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.
// 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)
$ ./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
| Type | Routing rule | Best 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.
// 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 found. After
UnbindQueue, publishes for that key fail instead of silently
disappearing.
DeleteExchange refuses to
remove an exchange that still has bindings. Unbind or delete its bound queues first.
$ ./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
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
| Field | Meaning |
|---|---|
ProducedCount | Messages published to the queue. |
AckedCount | Messages acknowledged by consumers. |
QueueDepth | Messages currently waiting in the ring buffer. |
InflightCount | Messages delivered but not yet acknowledged. |
RedeliveredCount | Messages requeued after an ack timeout. |
DlqCount | Messages dead-lettered after exhausting retries or expiring. |
ConsumerSessionCount | Connected consumers subscribed to the queue. |
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.
| Command | Description |
|---|---|
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. |
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
| Type | Sent by | Purpose |
|---|---|---|
publish | Producer | Route a message through an exchange to bound queues. |
subscribe | Consumer | Register a consumer session for a queue. |
ack | Consumer | Acknowledge a delivery by token. |
unsubscribe | Consumer | Leave the queue's round-robin rotation. |
metrics | Producer | Fetch broker-wide or per-queue metrics. |
create_queue | Producer | Create a named queue. |
delete_queue | Producer | Delete a named queue. |
list_queues | Producer | List all queue names. |
bind_queue | Producer | Bind a queue to an exchange with a binding key. |
unbind_queue | Producer | Remove a binding. |
create_exchange | Producer | Create an exchange with a type. |
delete_exchange | Producer | Delete an exchange without bindings. |
list_exchanges | Producer | List all exchange names. |
Example
{"type":"publish","exchange":"default","routing_key":"orders","message":{"payload":"order-1042 created"}}
{"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.
{"type":"bind_queue","exchange":"default","queue":"orders","binding_key":"orders"}
{"success":true,"error":""}
{
"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:
- Ack — the consumer processed it; it's removed permanently.
- Ack timeout — if the consumer doesn't ack within
AckTimeout, the message is requeued and redelivered. - Retry limit — after
MaxRetriesredeliveries, the message is moved to the queue's DLQ.
Default configuration
| Setting | Default | Effect |
|---|---|---|
AckTimeout | 100 s | How long a delivery may sit unacked before redelivery. |
RedeliveryInterval | 5 s | How often expired in-flight messages are swept. |
ExpiryInterval | 1 s | How often the expiry worker sweeps queues for TTL-expired messages. |
MaxRetries | 3 | Redeliveries before a message is dead-lettered. |
QueueSize | 20 | Capacity of each queue's ring buffer. |
ConsumerPrefetch | 10 | Max 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.