2022-01-05 18:44:49 +01:00
|
|
|
package jetstream
|
|
|
|
|
2022-02-02 14:32:48 +01:00
|
|
|
import (
|
|
|
|
"context"
|
2022-11-16 11:28:22 +01:00
|
|
|
"errors"
|
2022-02-02 14:32:48 +01:00
|
|
|
"fmt"
|
|
|
|
|
2022-03-07 17:40:56 +01:00
|
|
|
"github.com/getsentry/sentry-go"
|
2022-02-02 14:32:48 +01:00
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/sirupsen/logrus"
|
|
|
|
)
|
|
|
|
|
2022-08-31 13:21:56 +02:00
|
|
|
// JetStreamConsumer starts a durable consumer on the given subject with the
|
|
|
|
// given durable name. The function will be called when one or more messages
|
|
|
|
// is available, up to the maximum batch size specified. If the batch is set to
|
|
|
|
// 1 then messages will be delivered one at a time. If the function is called,
|
|
|
|
// the messages array is guaranteed to be at least 1 in size. Any provided NATS
|
|
|
|
// options will be passed through to the pull subscriber creation. The consumer
|
|
|
|
// will continue to run until the context expires, at which point it will stop.
|
2022-02-02 14:32:48 +01:00
|
|
|
func JetStreamConsumer(
|
2022-08-31 13:21:56 +02:00
|
|
|
ctx context.Context, js nats.JetStreamContext, subj, durable string, batch int,
|
|
|
|
f func(ctx context.Context, msgs []*nats.Msg) bool,
|
2022-02-02 14:32:48 +01:00
|
|
|
opts ...nats.SubOpt,
|
|
|
|
) error {
|
|
|
|
defer func() {
|
|
|
|
// If there are existing consumers from before they were pull
|
|
|
|
// consumers, we need to clean up the old push consumers. However,
|
|
|
|
// in order to not affect the interest-based policies, we need to
|
|
|
|
// do this *after* creating the new pull consumers, which have
|
|
|
|
// "Pull" suffixed to their name.
|
|
|
|
if _, err := js.ConsumerInfo(subj, durable); err == nil {
|
|
|
|
if err := js.DeleteConsumer(subj, durable); err != nil {
|
|
|
|
logrus.WithContext(ctx).Warnf("Failed to clean up old consumer %q", durable)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
|
|
|
name := durable + "Pull"
|
|
|
|
sub, err := js.PullSubscribe(subj, name, opts...)
|
|
|
|
if err != nil {
|
2022-03-07 17:40:56 +01:00
|
|
|
sentry.CaptureException(err)
|
2022-02-02 14:32:48 +01:00
|
|
|
return fmt.Errorf("nats.SubscribeSync: %w", err)
|
2022-01-05 18:44:49 +01:00
|
|
|
}
|
2022-02-02 14:32:48 +01:00
|
|
|
go func() {
|
|
|
|
for {
|
2022-04-27 16:29:49 +02:00
|
|
|
// If the parent context has given up then there's no point in
|
|
|
|
// carrying on doing anything, so stop the listener.
|
|
|
|
select {
|
|
|
|
case <-ctx.Done():
|
|
|
|
if err := sub.Unsubscribe(); err != nil {
|
|
|
|
logrus.WithContext(ctx).Warnf("Failed to unsubscribe %q", durable)
|
|
|
|
}
|
|
|
|
return
|
|
|
|
default:
|
|
|
|
}
|
2022-02-02 14:32:48 +01:00
|
|
|
// The context behaviour here is surprising — we supply a context
|
|
|
|
// so that we can interrupt the fetch if we want, but NATS will still
|
|
|
|
// enforce its own deadline (roughly 5 seconds by default). Therefore
|
|
|
|
// it is our responsibility to check whether our context expired or
|
|
|
|
// not when a context error is returned. Footguns. Footguns everywhere.
|
2022-08-31 13:21:56 +02:00
|
|
|
msgs, err := sub.Fetch(batch, nats.Context(ctx))
|
2022-02-02 14:32:48 +01:00
|
|
|
if err != nil {
|
|
|
|
if err == context.Canceled || err == context.DeadlineExceeded {
|
|
|
|
// Work out whether it was the JetStream context that expired
|
|
|
|
// or whether it was our supplied context.
|
|
|
|
select {
|
|
|
|
case <-ctx.Done():
|
|
|
|
// The supplied context expired, so we want to stop the
|
|
|
|
// consumer altogether.
|
|
|
|
return
|
|
|
|
default:
|
|
|
|
// The JetStream context expired, so the fetch probably
|
|
|
|
// just timed out and we should try again.
|
|
|
|
continue
|
|
|
|
}
|
2022-11-16 11:28:22 +01:00
|
|
|
} else if errors.Is(err, nats.ErrConsumerDeleted) {
|
|
|
|
// The consumer was deleted so stop.
|
|
|
|
return
|
2022-02-02 14:32:48 +01:00
|
|
|
} else {
|
|
|
|
// Something else went wrong, so we'll panic.
|
2022-03-07 17:40:56 +01:00
|
|
|
sentry.CaptureException(err)
|
2022-02-02 14:32:48 +01:00
|
|
|
logrus.WithContext(ctx).WithField("subject", subj).Fatal(err)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
if len(msgs) < 1 {
|
|
|
|
continue
|
|
|
|
}
|
2022-09-01 10:20:40 +02:00
|
|
|
for _, msg := range msgs {
|
|
|
|
if err = msg.InProgress(nats.Context(ctx)); err != nil {
|
|
|
|
logrus.WithContext(ctx).WithField("subject", subj).Warn(fmt.Errorf("msg.InProgress: %w", err))
|
|
|
|
sentry.CaptureException(err)
|
|
|
|
continue
|
|
|
|
}
|
2022-02-02 14:32:48 +01:00
|
|
|
}
|
2022-08-31 13:21:56 +02:00
|
|
|
if f(ctx, msgs) {
|
2022-09-01 10:20:40 +02:00
|
|
|
for _, msg := range msgs {
|
|
|
|
if err = msg.AckSync(nats.Context(ctx)); err != nil {
|
|
|
|
logrus.WithContext(ctx).WithField("subject", subj).Warn(fmt.Errorf("msg.AckSync: %w", err))
|
|
|
|
sentry.CaptureException(err)
|
|
|
|
}
|
2022-02-02 14:32:48 +01:00
|
|
|
}
|
|
|
|
} else {
|
2022-09-01 10:20:40 +02:00
|
|
|
for _, msg := range msgs {
|
|
|
|
if err = msg.Nak(nats.Context(ctx)); err != nil {
|
|
|
|
logrus.WithContext(ctx).WithField("subject", subj).Warn(fmt.Errorf("msg.Nak: %w", err))
|
|
|
|
sentry.CaptureException(err)
|
|
|
|
}
|
2022-02-02 14:32:48 +01:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
return nil
|
2022-01-05 18:44:49 +01:00
|
|
|
}
|