Skip to content

Bound a subscription and see what it sheds

What happens by default

Every subscription made through Client.Subscribe or Client.QueueSubscribe gets a bounded queue: 4,096 messages or 32 MiB, whichever comes first.

The underlying library's default is 500,000 messages. That is not a bound in any useful sense — it is a promise to exhaust memory at a moment nobody chose, long after the traffic that caused it has stopped being explicable.

gonats.ClientSettings{
    PendingMsgs:  16384,
    PendingBytes: 128 << 20,
}

The bound cannot be turned off

NATS reads a non-positive limit as unlimited, so PendingMsgs: -1 would hand back exactly the unbounded queue the bound exists to prevent. Non-positive limits are refused at Connect with ErrUnboundedSubscription.

An earlier version accepted them, and had a test asserting that was fine. It was not.

What happens when one fills

NATS discards from a full subscription queue. It does not delay the publisher.

That is the right behaviour, and worth being clear about why: a publisher held up by a slow subscriber stops draining whatever it is reading from, and the problem moves one layer up where it is harder to see. Shedding keeps the failure local to the consumer that caused it.

The cost is that shedding is invisible unless something counts it.

Seeing it

cli, err := client.Register(ctx, "nats-client", controller, settings, cred,
    client.InProcess(srv),
    client.WithLogger(log),
    client.OnShed(func(s client.Shed) {
        metrics.Shed.WithLabelValues(s.Subject, s.Queue).Set(float64(s.Dropped))
    }),
    client.OnError(func(subject string, err error) {
        metrics.AsyncErrors.WithLabelValues(subject).Inc()
    }),
)

Neither callback may block

OnShed and OnError run on the connection's single serial callback goroutine, ahead of the closed callback that Close waits for. A handler that blocks there starves the drain: every shutdown takes its full timeout, and the dispatcher goroutine never returns.

Send to a buffered channel or update a counter. Do not do I/O.

OnShed fires with a Shed carrying the subject, the queue group and the running count for that subscription. The count belongs to the subscriber that fell behind, because "something shed" answers a much less useful question than "which consumer is behind".

Subject is the subscription pattern, not the message subject — measured, and it is what makes this safe as a metric label. A subscription to tenant.*.events reports tenant.*.events however many tenants publish through it, so cardinality is bounded by the number of subscriptions and no tenant identifier ever reaches your metrics. Queue distinguishes two subscriptions to the same pattern in different groups, which the subject alone cannot.

The logger defaults to slog.Default(). A module whose stated purpose is that a shed is never silent cannot make reporting opt-in: the consumer most likely to lose messages unnoticed is the one that configured nothing.

OnError is not the same thing, and matters more than it looks

Some NATS failures never reach a return value. A publish the server later refuses on permissions returns nil from Publish and arrives at the asynchronous handler instead — so a consumer watching only returned errors sees a successful publish that never happened. Authentication failures, subscription limits and drain timeouts arrive the same way.

An earlier version of this module installed a handler that reported only slow-consumer errors and discarded the rest, which is worse than installing none at all: the library's own default at least prints them.

Why the handler form is required

Client.Subscribe takes a callback rather than returning a channel, and that is not a style preference.

SetPendingLimits returns ErrTypeSubscription for a channel subscription. The bound and the Dropped() count are only available on the handler form, so a channel-based API here would be one that cannot keep the promise this module exists to make.

Publishing through the module rather than around it

Client.Publish exists so a publish has one place to pass through — the place header propagation, metrics, deadlines and subject policy would go.

// Fire and forget: returns once the client has buffered it.
err := cli.Publish(ctx, "estate.events", payload)

// Ask to be told it arrived: the deadline is the request.
confirm, cancel := context.WithTimeout(ctx, 2*time.Second)
defer cancel()
err = cli.Publish(confirm, "estate.events", payload)

Confirmation costs a round trip, so it is asked for rather than assumed — the same reason there is no Len(). Without a deadline this behaves as NATS does, and a publish the server later refuses on permissions reaches OnError rather than the return value.