Summary
When a Sensor's JetStream pull subscription becomes permanently invalid while the underlying NATS connection stays open, the sensor stops delivering events and never recovers on its own. It spins in pullSubscribe indefinitely. Every external health signal stays green — the pod is Ready, the process is alive, and (see §3) the error is logged only once.
Recovery requires deleting the pod.
Environment
- argo-events
v1.9.11 (code below read at tag v1.9.11; both files are unchanged on master as of 2026-09-13)
- JetStream EventBus, 3 replicas
- Observed on
v1.9.3 as a ~21.5 hour outage; the code path is identical in v1.9.11
1. pullSubscribe treats a fatal error as retryable, forever
pkg/eventbus/jetstream/sensor/trigger_conn.go:196-225:
for {
msgs, fetchErr := subscription.Fetch(1, nats.MaxWait(time.Second*1))
if fetchErr != nil && !errors.Is(fetchErr, nats.ErrTimeout) {
if (previousErr != nil && previousErr.Error() != fetchErr.Error()) || time.Since(previousErrTime) > 10*time.Second {
conn.Logger.Errorf("failed to fetch messages for subscription %+v, %v, ...")
}
previousErr = fetchErr
previousErrTime = time.Now()
}
select {
case <-closeCh:
wg.Done()
return
default:
}
if fetchErr != nil && !errors.Is(fetchErr, nats.ErrTimeout) {
continue // <-- every fatal error lands here
}
...
}
Fetch errors are not classified. nats.ErrBadSubscription, nats.ErrSubscriptionClosed and nats: invalid subscription are permanent for this subscription object — retrying can never succeed — but they take the same continue as a transient one. The only exit from the loop is closeCh, which nothing signals in this scenario.
2. The existing reconnect daemon cannot see this failure
pkg/sensors/listener.go:319-349 already runs the recovery machinery — a 5s ticker that rebuilds the connection and resubscribes:
ticker := time.NewTicker(5 * time.Second)
...
case <-ticker.C:
if conn == nil || conn.IsClosed() {
triggerLogger.Info("EventBus connection lost, reconnecting...")
...
The guard is conn == nil || conn.IsClosed(). In this failure the connection is open and healthy — only the subscription is dead. So the daemon never fires, even though the code to fix the problem is sitting right there and is fully wired.
This is the root of it: the recovery path exists, it's just keyed on the wrong condition.
3. Why this is hard to observe (a secondary bug)
The log throttle at trigger_conn.go:200 is intended to print at most once per 10 seconds:
if (previousErr != nil && previousErr.Error() != fetchErr.Error()) || time.Since(previousErrTime) > 10*time.Second {
But line 205 sets previousErrTime = time.Now() on every erroring iteration. Inside a tight retry loop, time.Since(previousErrTime) is always ~0, so the > 10*time.Second branch can never be reached. The only remaining trigger is a changed error string — and errors.Join(ErrBadSubscription, ErrSubscriptionClosed) is constant.
Net effect: a stable wedge logs exactly once, at onset, and is then completely silent. Any log-based alert on this condition is structurally unable to fire on a sustained outage, which is precisely the case you most want to catch. A monitor built on the error line will look healthy for the entire incident.
4. Impact
- The sensor silently stops triggering. No restarts, no readiness failure, no crash — nothing a probe or a
kubectl get pods would surface.
- Duration is unbounded. Observed ~21.5 hours before manual intervention.
- On recovery via pod restart,
Initialize() deletes the durable consumer and the new PullSubscribe uses DeliverNew, so events that arrived during the outage window are not redelivered.
5. Suggested fix
Classify the error in pullSubscribe and return on the permanent ones (ErrBadSubscription, ErrSubscriptionClosed, invalid subscription), letting the existing listener.go daemon rebuild the subscription — rather than adding a parallel recovery path.
A naive version of this deadlocks, so if it's useful, four things need to hold:
wg.Done() must run on the new fatal-return path. Today it is reachable only via the closeCh branch (trigger_conn.go:210-214), so a wg.Wait() would hang forever on any other return.
processMsgsCloseCh and each pullSubscribeCloseCh need to be buffered to 1. shutdownSubscriptions does unbuffered sends to goroutines that may already have exited.
- The
msgChannel <- msg send at trigger_conn.go:223 should become a select with closeCh. The channel is unbuffered, and Go picks randomly among ready cases.
listener.go:339's closeSubCh <- struct{}{} should be a non-blocking select/default. The window between the surrounding statements can otherwise permanently kill the per-trigger daemon.
That's two files across two packages. Happy to open a PR with table-driven tests asserting no goroutine leak and no hang for each fatal error type if maintainers agree with the approach — I'd rather confirm the direction first than guess at it.
6. Separate request: please backport #4158 to release-1.9
#4158 ("fix(eventbus): close NATS connection when trigger connection init fails") merged to master on 2026-08-31, after the v1.9.11 release (2026-07-13), so it is not in any tagged release.
It addresses a connection leak on failed trigger-connection init, which presents as a steadily climbing CPU and goroutine count on the sensor — a distinct failure from the one above, and one that v1.9.11 users hit today with no released fix available. A release-1.9 backport would help anyone not able to run master.
Related
Summary
When a Sensor's JetStream pull subscription becomes permanently invalid while the underlying NATS connection stays open, the sensor stops delivering events and never recovers on its own. It spins in
pullSubscribeindefinitely. Every external health signal stays green — the pod is Ready, the process is alive, and (see §3) the error is logged only once.Recovery requires deleting the pod.
Environment
v1.9.11(code below read at tagv1.9.11; both files are unchanged onmasteras of 2026-09-13)v1.9.3as a ~21.5 hour outage; the code path is identical inv1.9.111.
pullSubscribetreats a fatal error as retryable, foreverpkg/eventbus/jetstream/sensor/trigger_conn.go:196-225:Fetcherrors are not classified.nats.ErrBadSubscription,nats.ErrSubscriptionClosedandnats: invalid subscriptionare permanent for this subscription object — retrying can never succeed — but they take the samecontinueas a transient one. The only exit from the loop iscloseCh, which nothing signals in this scenario.2. The existing reconnect daemon cannot see this failure
pkg/sensors/listener.go:319-349already runs the recovery machinery — a 5s ticker that rebuilds the connection and resubscribes:The guard is
conn == nil || conn.IsClosed(). In this failure the connection is open and healthy — only the subscription is dead. So the daemon never fires, even though the code to fix the problem is sitting right there and is fully wired.This is the root of it: the recovery path exists, it's just keyed on the wrong condition.
3. Why this is hard to observe (a secondary bug)
The log throttle at
trigger_conn.go:200is intended to print at most once per 10 seconds:But line 205 sets
previousErrTime = time.Now()on every erroring iteration. Inside a tight retry loop,time.Since(previousErrTime)is always ~0, so the> 10*time.Secondbranch can never be reached. The only remaining trigger is a changed error string — anderrors.Join(ErrBadSubscription, ErrSubscriptionClosed)is constant.Net effect: a stable wedge logs exactly once, at onset, and is then completely silent. Any log-based alert on this condition is structurally unable to fire on a sustained outage, which is precisely the case you most want to catch. A monitor built on the error line will look healthy for the entire incident.
4. Impact
kubectl get podswould surface.Initialize()deletes the durable consumer and the newPullSubscribeusesDeliverNew, so events that arrived during the outage window are not redelivered.5. Suggested fix
Classify the error in
pullSubscribeand return on the permanent ones (ErrBadSubscription,ErrSubscriptionClosed, invalid subscription), letting the existinglistener.godaemon rebuild the subscription — rather than adding a parallel recovery path.A naive version of this deadlocks, so if it's useful, four things need to hold:
wg.Done()must run on the new fatal-return path. Today it is reachable only via thecloseChbranch (trigger_conn.go:210-214), so awg.Wait()would hang forever on any other return.processMsgsCloseChand eachpullSubscribeCloseChneed to be buffered to 1.shutdownSubscriptionsdoes unbuffered sends to goroutines that may already have exited.msgChannel <- msgsend attrigger_conn.go:223should become aselectwithcloseCh. The channel is unbuffered, and Go picks randomly among ready cases.listener.go:339'scloseSubCh <- struct{}{}should be a non-blockingselect/default. The window between the surrounding statements can otherwise permanently kill the per-trigger daemon.That's two files across two packages. Happy to open a PR with table-driven tests asserting no goroutine leak and no hang for each fatal error type if maintainers agree with the approach — I'd rather confirm the direction first than guess at it.
6. Separate request: please backport #4158 to
release-1.9#4158 ("fix(eventbus): close NATS connection when trigger connection init fails") merged to
masteron 2026-08-31, after thev1.9.11release (2026-07-13), so it is not in any tagged release.It addresses a connection leak on failed trigger-connection init, which presents as a steadily climbing CPU and goroutine count on the sensor — a distinct failure from the one above, and one that
v1.9.11users hit today with no released fix available. Arelease-1.9backport would help anyone not able to runmaster.Related
not_planned)