diff --git a/go.mod b/go.mod index bde9f5b..eb8128b 100644 --- a/go.mod +++ b/go.mod @@ -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 ) @@ -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 diff --git a/go.sum b/go.sum index 12766e8..1ad0c00 100644 --- a/go.sum +++ b/go.sum @@ -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= @@ -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= @@ -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= diff --git a/internal/eventing/handler_image_nats.go b/internal/eventing/handler_image_nats.go index 3b4e8fa..0b5cf44 100644 --- a/internal/eventing/handler_image_nats.go +++ b/internal/eventing/handler_image_nats.go @@ -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. @@ -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 }