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.
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.