Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ require (
github.com/spf13/cobra v1.8.1
go.uber.org/zap v1.27.0
golang.org/x/net v0.29.0
golang.org/x/sync v0.13.0
gorm.io/driver/postgres v1.5.9
gorm.io/gorm v1.25.10
)
Expand Down Expand Up @@ -55,7 +56,6 @@ require (
golang.org/x/crypto v0.37.0 // indirect
golang.org/x/mod v0.19.0 // indirect
golang.org/x/oauth2 v0.18.0 // indirect
golang.org/x/sync v0.13.0 // indirect
golang.org/x/sys v0.32.0 // indirect
golang.org/x/text v0.24.0 // indirect
golang.org/x/tools v0.23.0 // indirect
Expand Down
8 changes: 8 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,8 @@ github.com/google/go-github/v39 v39.2.0 h1:rNNM311XtPOz5rDdsJXAp2o8F67X9FnROXTvt
github.com/google/go-github/v39 v39.2.0/go.mod h1:C1s8C5aCC9L+JXIYpJM5GYytdX52vC1bLvHEF1IhBrE=
github.com/google/go-querystring v1.1.0 h1:AnCroh3fv4ZBgVIf1Iwtovgjaw/GiKJo8M8yD/fhyJ8=
github.com/google/go-querystring v1.1.0/go.mod h1:Kcdr2DB4koayq7X8pmAG4sNG59So17icRSOU623lUBU=
github.com/google/go-tpm v0.9.3 h1:+yx0/anQuGzi+ssRqeD6WpXjW2L/V0dItUayO0i9sRc=
github.com/google/go-tpm v0.9.3/go.mod h1:h9jEsEECg7gtLis0upRBQU+GhYVH6jMjrFxI8u6bVUY=
github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
github.com/google/gofuzz v1.2.0 h1:xRy4A+RhZaiKjJ1bPfwQ8sedCA+YS2YcCHW6ec7JMi0=
github.com/google/gofuzz v1.2.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg=
Expand Down Expand Up @@ -239,6 +241,8 @@ github.com/mgutz/ansi v0.0.0-20170206155736-9520e82c474b h1:j7+1HpAFS1zy5+Q4qx1f
github.com/mgutz/ansi v0.0.0-20170206155736-9520e82c474b/go.mod h1:01TrycV0kFyexm33Z7vhZRXopbI8J3TDReVlkTgMUxE=
github.com/miekg/pkcs11 v1.1.1 h1:Ugu9pdy6vAYku5DEpVWVFPYnzV+bxB+iRdbuFSu7TvU=
github.com/miekg/pkcs11 v1.1.1/go.mod h1:XsNlhZGX73bx86s2hdc/FuaLm2CPZJemRLMA+WTFxgs=
github.com/minio/highwayhash v1.0.3 h1:kbnuUMoHYyVl7szWjSxJnxw11k2U709jqFPPmIUyD6Q=
github.com/minio/highwayhash v1.0.3/go.mod h1:GGYsuwP/fPD6Y9hMiXuapVvlIUEhFhMTh0rxU3ik1LQ=
github.com/minio/md5-simd v1.1.2 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34=
github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM=
github.com/minio/minio-go/v7 v7.0.67 h1:BeBvZWAS+kRJm1vGTMJYVjKUNoo0FoEt/wUWdUtfmh8=
Expand Down Expand Up @@ -280,6 +284,10 @@ github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
github.com/mxk/go-flowrate v0.0.0-20140419014527-cca7078d478f h1:y5//uYreIhSUg3J1GEMiLbxo1LJaP8RfCpH6pymGZus=
github.com/mxk/go-flowrate v0.0.0-20140419014527-cca7078d478f/go.mod h1:ZdcZmHo+o7JKHSa8/e818NopupXU1YMK5fe1lsApnBw=
github.com/nats-io/jwt/v2 v2.7.3 h1:6bNPK+FXgBeAqdj4cYQ0F8ViHRbi7woQLq4W29nUAzE=
github.com/nats-io/jwt/v2 v2.7.3/go.mod h1:GvkcbHhKquj3pkioy5put1wvPxs78UlZ7D/pY+BgZk4=
github.com/nats-io/nats-server/v2 v2.11.0 h1:fdwAT1d6DZW/4LUz5rkvQUe5leGEwjjOQYntzVRKvjE=
github.com/nats-io/nats-server/v2 v2.11.0/go.mod h1:leXySghbdtXSUmWem8K9McnJ6xbJOb0t9+NQ5HTRZjI=
github.com/nats-io/nats.go v1.45.0 h1:/wGPbnYXDM0pLKFjZTX+2JOw9TQPoIgTFrUaH97giwA=
github.com/nats-io/nats.go v1.45.0/go.mod h1:iRWIPokVIFbVijxuMQq4y9ttaBTMe0SFdlZfMDd+33g=
github.com/nats-io/nkeys v0.4.11 h1:q44qGV008kYd9W1b1nEBkNzvnWxtRSQ7A8BoqRrcfa0=
Expand Down
64 changes: 53 additions & 11 deletions internal/eventing/handler_image_nats.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,12 @@ import (
"github.com/weeb-vip/image-sync/internal/services/image_processor"
"github.com/weeb-vip/image-sync/internal/services/storage/minio"
"go.uber.org/zap"
"golang.org/x/sync/errgroup"
)

const (
maxRetries = 3
retryHeaderKey = "retry"
)

// EventingImageNats is the NATS counterpart of EventingImageKafka.
Expand Down Expand Up @@ -68,21 +74,57 @@ func EventingImageNats() error {

log.Info("Started NATS worker pool middleware", zap.Int("workers", cfg.WorkerConfig.KafkaImageProcessorWorkers), zap.Int("bufferSize", cfg.WorkerConfig.BufferSize))

processorInstance := processor.NewProcessor[*epNats.Message, image_processor.Payload](driver, cfg.NatsConfig.Subject, workerPoolMiddleware.Process)
retrySubject := cfg.NatsConfig.Subject + "-retry"
dlqSubject := cfg.NatsConfig.Subject + "-dlq"

log.Info("initializing backoff retry middleware", zap.String("subject", cfg.NatsConfig.Subject))
backoffRetryInstance := backoffretry.NewBackoffRetry[image_processor.Payload](driver, backoffretry.Config{
MaxRetries: 3,
HeaderKey: "retry",
RetryQueue: cfg.NatsConfig.Subject + "-retry",
// The retry consumer runs in this same process rather than a second
// deployment. It needs its own driver because the durable consumer name is
// driver-level configuration, not per-subject: two Consume calls on one
// driver would call CreateOrUpdateConsumer with the same durable name and
// different filter subjects, and the second would reconfigure the first.
retryDriver := epNats.NewNatsDriver(&epNats.Config{
URL: cfg.NatsConfig.URL,
ConsumerGroupName: cfg.NatsConfig.ConsumerGroupName + "-retry",
StreamName: cfg.NatsConfig.StreamName,
ConsumerAutoOffsetReset: &cfg.NatsConfig.Offset,
})
defer func(d drivers.Driver[*epNats.Message]) {
if err := d.Close(); err != nil {
log.Error("Error closing NATS retry driver", zap.String("error", err.Error()))
}
}(retryDriver)

processorInstance := processor.NewProcessor[*epNats.Message, image_processor.Payload](driver, cfg.NatsConfig.Subject, workerPoolMiddleware.Process).
AddMiddleware(backoffretry.NewBackoffRetry[image_processor.Payload](driver, backoffretry.Config{
MaxRetries: maxRetries,
HeaderKey: retryHeaderKey,
RetryQueue: retrySubject,
}).Process)

// Exhausted retries go to a dead-letter subject rather than back onto the
// retry subject. ep acks and drops a message once the counter reaches
// MaxRetries, so cycling it here would make a permanently failing message
// disappear with no record.
retryProcessorInstance := processor.NewProcessor[*epNats.Message, image_processor.Payload](retryDriver, retrySubject, workerPoolMiddleware.Process).
AddMiddleware(backoffretry.NewBackoffRetry[image_processor.Payload](retryDriver, backoffretry.Config{
MaxRetries: maxRetries,
HeaderKey: retryHeaderKey,
RetryQueue: dlqSubject,
}).Process)

log.Info("Starting NATS processors",
zap.String("subject", cfg.NatsConfig.Subject),
zap.String("retry_subject", retrySubject),
zap.String("dlq_subject", dlqSubject))

log.Info("Starting NATS processor", zap.String("subject", cfg.NatsConfig.Subject))
err := processorInstance.
AddMiddleware(backoffRetryInstance.Process).
Run(ctx)
// One consumer returning must stop the other: Consume blocks until its
// iterator is stopped, so without cancelling here a dead main consumer
// would leave the process alive and apparently healthy.
group, groupCtx := errgroup.WithContext(ctx)
group.Go(func() error { return processorInstance.Run(groupCtx) })
group.Go(func() error { return retryProcessorInstance.Run(groupCtx) })

if err != nil && ctx.Err() == nil { // Ignore error if caused by context cancellation
if err := group.Wait(); err != nil && ctx.Err() == nil { // Ignore error if caused by context cancellation
log.Error("Error consuming messages", zap.String("error", err.Error()))
return err
}
Expand Down
Loading