diff --git a/cdc/init.go b/cdc/init.go index 4ddf22ac..7ed9a9a8 100644 --- a/cdc/init.go +++ b/cdc/init.go @@ -2,10 +2,14 @@ package cdc import ( "context" + "errors" + "log" + + "go.uber.org/fx" + "go.uber.org/zap" + "github.com/tidepool-org/go-common/asyncevents" "github.com/tidepool-org/go-common/events" - "go.uber.org/fx" - "log" ) func AttachConsumerGroupHooks(cg events.EventConsumer, lifecycle fx.Lifecycle, shutdowner fx.Shutdowner) { @@ -27,25 +31,75 @@ func AttachConsumerGroupHooks(cg events.EventConsumer, lifecycle fx.Lifecycle, s }) } -func AttachSaramaRunnerHooks(runner asyncevents.SaramaEventsRunner, lifecycle fx.Lifecycle, shutdowner fx.Shutdowner) { - adapted := asyncevents.NewSaramaRunner(runner) - lifecycle.Append(fx.Hook{ - OnStart: func(ctx context.Context) error { - go func() { - if err := adapted.Initialize(); err != nil { - log.Printf("unable to initialize runner: %v", err) - } - if err := adapted.Run(); err != nil { - log.Printf("error from runner: %v", err) - if err := shutdowner.Shutdown(); err != nil { - log.Printf("error shutting down runner: %v", err) - } - } - }() - return nil - }, - OnStop: func(ctx context.Context) error { - return adapted.Terminate() - }, - }) +type SaramaFxAdapterConfig struct { + Consumers []asyncevents.Runner `group:"runners"` + Logger *zap.SugaredLogger + Shutdowner fx.Shutdowner +} + +// SaramaFxAdapter adds fx.Lifecycle and Shutdown support to go-common's +// asyncevents.NonBlockingStartStopAdapter. +type SaramaFxAdapter struct { + *SaramaFxAdapterConfig + + adapters []*asyncevents.NonBlockingStartStopAdapter +} + +type SaramaFxAdapterParams struct { + fx.In + + SaramaFxAdapterConfig +} + +func NewSaramaFxAdapter(p SaramaFxAdapterParams) *SaramaFxAdapter { + return &SaramaFxAdapter{ + SaramaFxAdapterConfig: &p.SaramaFxAdapterConfig, + } } + +// Start implements fx.HookFunc. +func (f *SaramaFxAdapter) Start(ctx context.Context) error { + adapters := []*asyncevents.NonBlockingStartStopAdapter{} + for _, consumer := range f.Consumers { + adapter := asyncevents.NewNonBlockingStartStopAdapter(consumer, f.onError) + if err := adapter.Start(ctx); err != nil { + return err + } + adapters = append(adapters, adapter) + } + f.adapters = adapters + return nil +} + +// onError is passed to asyncevents.NonBlockingStartStopAdapter to be called on an error. +// +// Using an fx.Shutdowner, we can then shutdown the service. +func (f *SaramaFxAdapter) onError(err error) { + f.Logger.With("error", err).Info("consumer exited unexpectedly") + if err := f.Shutdowner.Shutdown(fx.ExitCode(1)); err != nil { + f.Logger.With("error", err).Info("shutting down fx") + } +} + +// Start implements fx.HookFunc. +func (f *SaramaFxAdapter) Stop(ctx context.Context) error { + var joinedErrs error + for _, adapter := range f.adapters { + if err := adapter.Stop(ctx); err != nil { + joinedErrs = errors.Join(joinedErrs, err) + } + } + return joinedErrs +} + +var Module = fx.Provide( + fx.Annotate( + NewSaramaFxAdapter, + fx.OnStart(func(ctx context.Context, sfa *SaramaFxAdapter) error { + return sfa.Start(ctx) + }), + fx.OnStop(func(ctx context.Context, sfa *SaramaFxAdapter) error { + return sfa.Stop(ctx) + }), + ), +) diff --git a/go.mod b/go.mod index 496f769c..5c1ea109 100644 --- a/go.mod +++ b/go.mod @@ -15,7 +15,7 @@ require ( github.com/onsi/gomega v1.36.2 github.com/tidepool-org/clinic/client v0.0.0-20250521172904-b61821ac9973 github.com/tidepool-org/clinic/redox_models v0.0.0-20250521172904-b61821ac9973 - github.com/tidepool-org/go-common v0.12.3-0.20250613120630-656deb326ad3 + github.com/tidepool-org/go-common v0.12.3-0.20250625233514-f3a3357c3d91 github.com/tidepool-org/hydrophone/client v0.0.0-20250317164837-a8cd51fd6677 go.mongodb.org/mongo-driver v1.17.3 go.uber.org/fx v1.23.0 diff --git a/go.sum b/go.sum index 19017301..f7a20fce 100644 --- a/go.sum +++ b/go.sum @@ -107,8 +107,8 @@ github.com/tidepool-org/clinic/client v0.0.0-20250521172904-b61821ac9973 h1:KbPq github.com/tidepool-org/clinic/client v0.0.0-20250521172904-b61821ac9973/go.mod h1:r4oWW+WA1IzIcB7Y4q9iun2rcihTVXVVYmB4F9dq6HA= github.com/tidepool-org/clinic/redox_models v0.0.0-20250521172904-b61821ac9973 h1:L8bcUPW2VIkN96M1P/PKPaDvnc0j5CXL1I2gxn6nF7A= github.com/tidepool-org/clinic/redox_models v0.0.0-20250521172904-b61821ac9973/go.mod h1:bQ9DZxk015RhmGG1tR6jRScP9KxyHvS8tzPbVtr82DE= -github.com/tidepool-org/go-common v0.12.3-0.20250613120630-656deb326ad3 h1:X4+QkiivmVvkFi8B4GCZryvso4JLoPKkKTfC5nqtT44= -github.com/tidepool-org/go-common v0.12.3-0.20250613120630-656deb326ad3/go.mod h1:v93bMGDHiHcltQY5s7LYTTEe3u9CiGWqBFKah6C0650= +github.com/tidepool-org/go-common v0.12.3-0.20250625233514-f3a3357c3d91 h1:zWimBR1JVcUwK+n9DdNNw1HWydTzHQCH3LxkckoLbEA= +github.com/tidepool-org/go-common v0.12.3-0.20250625233514-f3a3357c3d91/go.mod h1:v93bMGDHiHcltQY5s7LYTTEe3u9CiGWqBFKah6C0650= github.com/tidepool-org/hydrophone/client v0.0.0-20250317164837-a8cd51fd6677 h1:P3C1YTvLHu7NFHOeh6NDPo5ieVXcF9a+SY4FTToR8B4= github.com/tidepool-org/hydrophone/client v0.0.0-20250317164837-a8cd51fd6677/go.mod h1:gon+x+jAh8DZZ2hD23fBWqrYwOizVSwIBbxEsuXCbZ4= github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= diff --git a/redox/consumer_message.go b/redox/consumer_message.go index 5d8f8cfa..ebe31c39 100644 --- a/redox/consumer_message.go +++ b/redox/consumer_message.go @@ -3,14 +3,16 @@ package redox import ( "context" "fmt" + "time" + "github.com/IBM/sarama" - "github.com/tidepool-org/clinic-worker/cdc" - models "github.com/tidepool-org/clinic/redox_models" - "github.com/tidepool-org/go-common/asyncevents" "go.mongodb.org/mongo-driver/bson" "go.uber.org/fx" "go.uber.org/zap" - "time" + + "github.com/tidepool-org/clinic-worker/cdc" + models "github.com/tidepool-org/clinic/redox_models" + "github.com/tidepool-org/go-common/asyncevents" ) const ( @@ -33,7 +35,7 @@ type MessageCDCConsumer struct { orderProcessor NewOrderProcessor } -func NewMessageCDCConsumer(p MessageCDCConsumerParams) (asyncevents.SaramaEventsRunner, error) { +func NewMessageCDCConsumer(p MessageCDCConsumerParams) (asyncevents.Runner, error) { if !p.Config.Enabled { return &cdc.DisabledSaramaEventsRunner{}, nil } @@ -48,25 +50,27 @@ func NewMessageCDCConsumer(p MessageCDCConsumerParams) (asyncevents.SaramaEvents prefixedTopics := []string{config.GetPrefixedTopic()} - runnerCfg := asyncevents.SaramaRunnerConfig{ - Brokers: config.KafkaBrokers, - GroupID: config.KafkaConsumerGroup, - Topics: prefixedTopics, - Sarama: config.SaramaConfig, - MessageConsumer: &MessageCDCConsumer{ - config: p.Config, - logger: p.Logger, - orderProcessor: p.OrderProcessor, - }, - } - delays := []time.Duration{0, time.Second * 60, time.Second * 300} logger := &cdc.AsynceventsLoggerAdapter{ SugaredLogger: p.Logger, } - eventsRunner := asyncevents.NewCascadingSaramaEventsRunner(runnerCfg, logger, delays, defaultTimeout) - return eventsRunner, nil + managerConfig := asyncevents.CascadingSaramaEventsManagerConfig{ + Consumer: &MessageCDCConsumer{ + config: p.Config, + logger: p.Logger, + orderProcessor: p.OrderProcessor, + }, + Brokers: config.KafkaBrokers, + GroupID: config.KafkaConsumerGroup, + Topics: prefixedTopics, + ConsumptionTimeout: defaultTimeout, + Delays: delays, + Logger: logger, + Sarama: config.SaramaConfig, + } + eventsManager := asyncevents.NewCascadingSaramaEventsManager(managerConfig) + return eventsManager, nil } func (c *MessageCDCConsumer) Consume(ctx context.Context, session sarama.ConsumerGroupSession, msg *sarama.ConsumerMessage) error { diff --git a/redox/consumer_scheduled.go b/redox/consumer_scheduled.go index 5ff8a5cc..c942f9a1 100644 --- a/redox/consumer_scheduled.go +++ b/redox/consumer_scheduled.go @@ -3,13 +3,15 @@ package redox import ( "context" "fmt" + "time" + "github.com/IBM/sarama" - "github.com/tidepool-org/clinic-worker/cdc" - "github.com/tidepool-org/go-common/asyncevents" "go.mongodb.org/mongo-driver/bson" "go.uber.org/fx" "go.uber.org/zap" - "time" + + "github.com/tidepool-org/clinic-worker/cdc" + "github.com/tidepool-org/go-common/asyncevents" ) const ( @@ -25,7 +27,7 @@ type ScheduledSummaryAndReportsCDCConsumerParams struct { Processor ScheduledSummaryAndReportProcessor } -func NewScheduledSummaryAndReportsRunner(p ScheduledSummaryAndReportsCDCConsumerParams) (asyncevents.SaramaEventsRunner, error) { +func NewScheduledSummaryAndReportsRunner(p ScheduledSummaryAndReportsCDCConsumerParams) (asyncevents.Runner, error) { if !p.Config.Enabled { return &cdc.DisabledSaramaEventsRunner{}, nil } @@ -40,25 +42,27 @@ func NewScheduledSummaryAndReportsRunner(p ScheduledSummaryAndReportsCDCConsumer prefixedTopics := []string{config.GetPrefixedTopic()} - runnerCfg := asyncevents.SaramaRunnerConfig{ - Brokers: config.KafkaBrokers, - GroupID: config.KafkaConsumerGroup, - Topics: prefixedTopics, - Sarama: config.SaramaConfig, - MessageConsumer: &ScheduledSummaryAndReportsCDCConsumer{ - Config: p.Config, - Logger: p.Logger, - Processor: p.Processor, - }, - } - delays := []time.Duration{0, time.Second * 60, time.Second * 300} logger := &cdc.AsynceventsLoggerAdapter{ SugaredLogger: p.Logger, } - eventsRunner := asyncevents.NewCascadingSaramaEventsRunner(runnerCfg, logger, delays, defaultTimeout) - return eventsRunner, nil + managerConfig := asyncevents.CascadingSaramaEventsManagerConfig{ + Consumer: &ScheduledSummaryAndReportsCDCConsumer{ + Config: p.Config, + Logger: p.Logger, + Processor: p.Processor, + }, + Brokers: config.KafkaBrokers, + GroupID: config.KafkaConsumerGroup, + Topics: prefixedTopics, + ConsumptionTimeout: defaultTimeout, + Delays: delays, + Logger: logger, + Sarama: config.SaramaConfig, + } + eventsManager := asyncevents.NewCascadingSaramaEventsManager(managerConfig) + return eventsManager, nil } // ScheduledSummaryAndReportsCDCConsumer is kafka consumer for scheduled summary and reports CDC events diff --git a/redox/module.go b/redox/module.go index 2eecd3e2..bc4bee49 100644 --- a/redox/module.go +++ b/redox/module.go @@ -1,10 +1,12 @@ package redox import ( + "time" + "github.com/kelseyhightower/envconfig" - "github.com/tidepool-org/clinic-worker/report" "go.uber.org/fx" - "time" + + "github.com/tidepool-org/clinic-worker/report" ) var Module = fx.Provide( diff --git a/vendor/github.com/tidepool-org/go-common/asyncevents/cascade.go b/vendor/github.com/tidepool-org/go-common/asyncevents/cascade.go index 0fb8196f..cdd5e9e8 100644 --- a/vendor/github.com/tidepool-org/go-common/asyncevents/cascade.go +++ b/vendor/github.com/tidepool-org/go-common/asyncevents/cascade.go @@ -5,175 +5,198 @@ import ( "context" "errors" "fmt" - "github.com/IBM/sarama" "log/slog" "os" "strconv" "sync" "time" + + "github.com/IBM/sarama" ) -// SaramaRunner interfaces between [events.Runner] and go-common's -// [SaramaEventsConsumer]. +// CascadingSaramaMessageConsumer cascades messages that failed to be consumed to another +// topic. It is an implementation of [SaramaMessageConsumer]. // -// This means providing Initialize(), Run(), and Terminate() to satisfy events.Runner, while -// under the hood calling SaramaEventConsumer's Run(), and canceling its Context as -// appropriate. -type SaramaRunner struct { - eventsRunner SaramaEventsRunner - cancelCtx context.CancelFunc - cancelMu sync.Mutex +// It also sets an adjustable delay via the "not-before" and "failures" headers so that as +// the message moves from topic to topic, the time between processing is increased according +// to [FailuresToDelay]. +type CascadingSaramaMessageConsumer struct { + Consumer SaramaMessageConsumer + NextTopic string + Producer LimitedAsyncProducer + Logger Logger } -func NewSaramaRunner(eventsRunner SaramaEventsRunner) *SaramaRunner { - return &SaramaRunner{ - eventsRunner: eventsRunner, - } -} +// Consume implements [SaramaMessageConsumer]. +func (c *CascadingSaramaMessageConsumer) Consume(ctx context.Context, + session sarama.ConsumerGroupSession, msg *sarama.ConsumerMessage) (err error) { -// SaramaEventsRunner is implemented by go-common's [SaramaEventsRunner]. -type SaramaEventsRunner interface { - Run(ctx context.Context) error + if err := c.Consumer.Consume(ctx, session, msg); err != nil { + txnErr := c.withTxn(ctx, func() error { + select { + case <-ctx.Done(): + if ctxErr := ctx.Err(); !errors.Is(ctxErr, context.Canceled) { + return ctxErr + } + return nil + case c.Producer.Input() <- c.cascadeMessage(ctx, msg): + c.Logger.Log(ctx, slog.LevelInfo, "cascaded", "from", msg.Topic, "to", c.NextTopic) + return nil + } + }) + if txnErr != nil { + c.Logger.Log(ctx, slog.LevelInfo, "Unable to complete cascading transaction", "error", err) + return err + } + } + return nil } -// SaramaRunnerConfig collects values needed to initialize a SaramaRunner. -// -// This provides isolation for the SaramaRunner from ConfigReporter, -// envconfig, or any of the other options in platform for reading config -// values. -type SaramaRunnerConfig struct { - Brokers []string - GroupID string - Topics []string - MessageConsumer SaramaMessageConsumer - - Sarama *sarama.Config +// withTxn wraps a function with a transaction that is aborted if an error is returned. +func (c *CascadingSaramaMessageConsumer) withTxn(ctx context.Context, f func() error) (err error) { + if err := c.Producer.BeginTxn(); err != nil { + return fmt.Errorf("unable to begin transaction: %w", err) + } + defer func(err *error) { + if err != nil && *err != nil { + if abortErr := c.Producer.AbortTxn(); abortErr != nil { + c.Logger.Log(ctx, slog.LevelInfo, "Unable to abort transaction", "error", abortErr) + } + return + } + if commitErr := c.Producer.CommitTxn(); commitErr != nil { + c.Logger.Log(ctx, slog.LevelInfo, "Unable to commit transaction", "error", commitErr) + } + }(&err) + return f() } -func (r *SaramaRunner) Initialize() error { return nil } - -// Run adapts platform's event.Runner to work with go-common's -// SaramaEventsConsumer. -func (r *SaramaRunner) Run() error { - if r.eventsRunner == nil { - return errors.New("unable to run SaramaRunner, eventsRunner is nil") - } +// cascadeMessage to the next topic. +func (c *CascadingSaramaMessageConsumer) cascadeMessage(ctx context.Context, + msg *sarama.ConsumerMessage) *sarama.ProducerMessage { - r.cancelMu.Lock() - ctx, err := func() (context.Context, error) { - defer r.cancelMu.Unlock() - if r.cancelCtx != nil { - return nil, errors.New("unable to Run SaramaRunner, it's already initialized") - } - var ctx context.Context - ctx, r.cancelCtx = context.WithCancel(context.Background()) - return ctx, nil - }() - if err != nil { - return err + pHeaders := make([]sarama.RecordHeader, len(msg.Headers)) + for idx, header := range msg.Headers { + pHeaders[idx] = *header } - if err := r.eventsRunner.Run(ctx); err != nil { - return fmt.Errorf("unable to Run SaramaRunner: %w", err) + return &sarama.ProducerMessage{ + Key: sarama.ByteEncoder(msg.Key), + Value: sarama.ByteEncoder(msg.Value), + Topic: c.NextTopic, + Headers: c.updateCascadeHeaders(ctx, pHeaders), } - return nil } -// Terminate adapts platform's event.Runner to work with go-common's -// SaramaEventsConsumer. -func (r *SaramaRunner) Terminate() error { - r.cancelMu.Lock() - defer r.cancelMu.Unlock() - if r.cancelCtx == nil { - return errors.New("unable to Terminate SaramaRunner, it's not running") - } - r.cancelCtx() - return nil -} +// updateCascadeHeaders calculates not before and failures header values. +// +// Existing not before and failures headers will be dropped in place of the new ones. +func (c *CascadingSaramaMessageConsumer) updateCascadeHeaders(ctx context.Context, + headers []sarama.RecordHeader) []sarama.RecordHeader { -// CappedExponentialBinaryDelay builds delay functions that use exponential -// binary backoff with a maximum duration. -func CappedExponentialBinaryDelay(cap time.Duration) func(int) time.Duration { - return func(tries int) time.Duration { - b := DelayExponentialBinary(tries) - if b > cap { - return cap + failures := 0 + notBefore := time.Now() + + keep := make([]sarama.RecordHeader, 0, len(headers)) + for _, header := range headers { + switch { + case bytes.Equal(header.Key, HeaderNotBefore): + continue // Drop this header, we'll add a new version below. + case bytes.Equal(header.Key, HeaderFailures): + parsed, err := strconv.ParseInt(string(header.Value), 10, 32) + if err != nil { + c.Logger.Log(ctx, slog.LevelInfo, "Unable to parse consumption failures count", "error", err) + } else { + failures = int(parsed) + notBefore = notBefore.Add(FailuresToDelay[failures]) + } + continue // Drop this header, we'll add a new version below. } - return b + keep = append(keep, header) } + + keep = append(keep, sarama.RecordHeader{ + Key: HeaderNotBefore, + Value: []byte(notBefore.Format(NotBeforeTimeFormat)), + }) + keep = append(keep, sarama.RecordHeader{ + Key: HeaderFailures, + Value: []byte(strconv.Itoa(failures + 1)), + }) + + return keep } -// CascadingSaramaEventsRunner manages multiple sarama consumer groups to execute a -// topic-cascading retry process. -// -// The topic names are generated from Config.Topics combined with Delays. If given a single -// topic "updates", and delays: 0s, 1s, and 5s, then the following topics will be consumed: -// updates, updates-retry-1s, updates-retry-5s. The consumer of the updates-retry-5s topic -// will write failed messages to updates-dead. -// -// The inspiration for this system was drawn from -// https://www.uber.com/blog/reliable-reprocessing/ -type CascadingSaramaEventsRunner struct { - Config SaramaRunnerConfig +// CascadingSaramaEventsManagerConfig for a [CascadingSaramaEventsManager]. +type CascadingSaramaEventsManagerConfig struct { + Consumer SaramaMessageConsumer + + Brokers []string + GroupID string + Topics []string ConsumptionTimeout time.Duration Delays []time.Duration Logger Logger SaramaBuilders SaramaBuilders + Sarama *sarama.Config } -func NewCascadingSaramaEventsRunner(config SaramaRunnerConfig, logger Logger, - delays []time.Duration, consumptionTimeout time.Duration) *CascadingSaramaEventsRunner { - - return &CascadingSaramaEventsRunner{ - Config: config, - Delays: delays, - Logger: logger, - SaramaBuilders: DefaultSaramaBuilders{}, - ConsumptionTimeout: consumptionTimeout, - } +// CascadingSaramaEventsManager manages multiple Sarama consumer groups to execute a +// topic-cascading retry process. It coordinates multiple [SaramaConsumerGroupManager] +// instances to achieve this. +// +// The topics' names are generated from a combination of the configured topics and +// configured delays. For example, if configured with a topic "updates", and delays: 0s, 1s, +// and 5s, then the following topics will be consumed: updates, updates-retry-1s, +// updates-retry-5s. The consumer of the updates-retry-5s topic will write failed messages +// to updates-dead. +// +// The inspiration for this system was drawn from +// https://www.uber.com/blog/reliable-reprocessing/ +type CascadingSaramaEventsManager struct { + CascadingSaramaEventsManagerConfig } -// LimitedAsyncProducer restricts the [sarama.AsyncProducer] interface to ensure that its -// recipient isn't able to call Close(), thereby opening the potential for a panic when -// writing to a closed channel. -type LimitedAsyncProducer interface { - AbortTxn() error - BeginTxn() error - CommitTxn() error - Input() chan<- *sarama.ProducerMessage +func NewCascadingSaramaEventsManager(config CascadingSaramaEventsManagerConfig) *CascadingSaramaEventsManager { + if config.SaramaBuilders == nil { + config.SaramaBuilders = &DefaultSaramaBuilders{} + } + return &CascadingSaramaEventsManager{ + CascadingSaramaEventsManagerConfig: config, + } } -func (r *CascadingSaramaEventsRunner) Run(ctx context.Context) error { - if len(r.Config.Topics) == 0 { +func (c *CascadingSaramaEventsManager) Run(ctx context.Context) error { + if len(c.Topics) == 0 { return errors.New("no topics") } - if len(r.Delays) == 0 { + if len(c.Delays) == 0 { return errors.New("no delays") } producersCtx, cancel := context.WithCancel(ctx) defer cancel() var wg sync.WaitGroup - errs := make(chan error, len(r.Config.Topics)*len(r.Delays)) + errs := make(chan error, len(c.Topics)*len(c.Delays)) defer func() { - r.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsRunner: waiting for consumers") + c.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsManager: waiting for managers") wg.Wait() - r.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsRunner: all consumers returned") + c.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsManager: all managers returned") close(errs) }() - for _, topic := range r.Config.Topics { - for idx, delay := range r.Delays { - producerCfg := r.producerConfig(idx, delay) - // The producer is built here rather than in buildConsumer() to control when - // producer is closed. Were the producer to be closed before consumer.Run() - // returns, it would be possible for consumer to write to the producer's + for _, topic := range c.Topics { + for idx, delay := range c.Delays { + producerCfg := c.producerConfig(idx, delay) + // The producer is built here rather than in buildManager() to control when + // producer is closed. Were the producer to be closed before manager.Run() + // returns, it would be possible for manager to write to the producer's // Inputs() channel, which if closed, would cause a panic. - producer, err := r.SaramaBuilders.NewAsyncProducer(r.Config.Brokers, producerCfg) + producer, err := c.SaramaBuilders.NewAsyncProducer(c.Brokers, producerCfg) if err != nil { - return fmt.Errorf("unable to build async producer %s: %w", r.Config.GroupID, err) + return fmt.Errorf("unable to build async producer %s: %w", c.GroupID, err) } - consumer, err := r.buildConsumer(producersCtx, idx, producer, delay, topic) + manager, err := c.buildManager(producersCtx, idx, producer, delay, topic) if err != nil { return err } @@ -183,35 +206,35 @@ func (r *CascadingSaramaEventsRunner) Run(ctx context.Context) error { defer func() { closeErr := producer.Close() if closeErr != nil { - r.Logger.Log(producersCtx, slog.LevelInfo, "CascadingSaramaEventsRunner: unable to close producer", "error", closeErr) + c.Logger.Log(producersCtx, slog.LevelInfo, "CascadingSaramaEventsManager: unable to close producer", "error", closeErr) } wg.Done() }() - if err := consumer.Run(producersCtx); err != nil { + if err := manager.Run(producersCtx); err != nil { errs <- fmt.Errorf("topics[%q]: %s", topic, err) } - r.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsRunner: consumer go proc returning", "topic", topic) + c.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsManager: manager go proc returning", "topic", topic) }(topic) } } select { case <-ctx.Done(): - r.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsRunner: context is done") + c.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsManager: context is done") return nil case err := <-errs: - r.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsRunner: Run(): error from consumer", "error", err) + c.Logger.Log(ctx, slog.LevelDebug, "CascadingSaramaEventsManager: Run(): error from manager", "error", err) return err } } -func (r *CascadingSaramaEventsRunner) producerConfig(idx int, delay time.Duration) *sarama.Config { - uniqueConfig := *r.Config.Sarama +func (c *CascadingSaramaEventsManager) producerConfig(idx int, delay time.Duration) *sarama.Config { + uniqueConfig := *c.Sarama hostID := os.Getenv("HOSTNAME") // set by default in kubernetes pods if hostID == "" { hostID = fmt.Sprintf("%d-%d", time.Now().UnixNano()/int64(time.Second), os.Getpid()) } - txnID := fmt.Sprintf("%s-%s-%d-%s", r.Config.GroupID, delay.String(), idx, hostID) + txnID := fmt.Sprintf("%s-%s-%d-%s", c.GroupID, delay.String(), idx, hostID) uniqueConfig.Producer.Transaction.ID = txnID uniqueConfig.Producer.Idempotent = true uniqueConfig.Producer.RequiredAcks = sarama.WaitForAll @@ -241,48 +264,50 @@ func (DefaultSaramaBuilders) NewConsumerGroup(brokers []string, groupID string, return sarama.NewConsumerGroup(brokers, groupID, config) } -func (r *CascadingSaramaEventsRunner) buildConsumer(ctx context.Context, idx int, +// buildManager returns a [SaramaConsumerGroupManager] that manages multiple, composed, +// [SaramaMessageConsumer]s according to its topics and delays configuration. +func (c *CascadingSaramaEventsManager) buildManager(ctx context.Context, idx int, producer LimitedAsyncProducer, delay time.Duration, baseTopic string) ( - *SaramaEventsConsumer, error) { + *SaramaConsumerGroupManager, error) { - groupID := r.Config.GroupID + groupID := c.GroupID if delay > 0 { groupID += "-retry-" + delay.String() } - group, err := r.SaramaBuilders.NewConsumerGroup(r.Config.Brokers, groupID, - r.Config.Sarama) + group, err := c.SaramaBuilders.NewConsumerGroup(c.Brokers, groupID, + c.Sarama) if err != nil { return nil, fmt.Errorf("unable to build sarama consumer group %s: %w", groupID, err) } - var consumer = r.Config.MessageConsumer - if len(r.Delays) > 0 { + var consumer SaramaMessageConsumer = c.Consumer + if len(c.Delays) > 0 { nextTopic := baseTopic + "-dead" - if idx+1 < len(r.Delays) { - nextTopic = baseTopic + "-retry-" + r.Delays[idx+1].String() + if idx+1 < len(c.Delays) { + nextTopic = baseTopic + "-retry-" + c.Delays[idx+1].String() } - consumer = &CascadingConsumer{ + consumer = &CascadingSaramaMessageConsumer{ Consumer: consumer, NextTopic: nextTopic, Producer: producer, - Logger: r.Logger, + Logger: c.Logger, } } if delay > 0 { consumer = &NotBeforeConsumer{ Consumer: consumer, - Logger: r.Logger, + Logger: c.Logger, } } - handler := NewSaramaConsumerGroupHandler(r.Logger, consumer, r.ConsumptionTimeout) + handler := NewSaramaConsumerGroupHandler(c.Logger, consumer, c.ConsumptionTimeout) topic := baseTopic if delay > 0 { topic += "-retry-" + delay.String() } - r.Logger.Log(ctx, slog.LevelDebug, "creating consumer", "topic", topic) + c.Logger.Log(ctx, slog.LevelDebug, "creating consumer", "topic", topic) - return NewSaramaEventsConsumer(group, handler, topic), nil + return NewSaramaConsumerGroupManager(group, handler, topic), nil } // NotBeforeConsumer delays consumption until a specified time. @@ -351,108 +376,24 @@ func (c *NotBeforeConsumer) notBeforeFromMsgHeaders(msg *sarama.ConsumerMessage) return time.Time{}, fmt.Errorf("header not found: x-tidepool-not-before") } -// CascadingConsumer cascades messages that failed to be consumed to another topic. -// -// It also sets an adjustable delay via the "not-before" and "failures" headers so that as -// the message moves from topic to topic, the time between processing is increased according -// to [FailuresToDelay]. -type CascadingConsumer struct { - Consumer SaramaMessageConsumer - NextTopic string - Producer LimitedAsyncProducer - Logger Logger -} - -func (c *CascadingConsumer) Consume(ctx context.Context, session sarama.ConsumerGroupSession, - msg *sarama.ConsumerMessage) (err error) { - - if err := c.Consumer.Consume(ctx, session, msg); err != nil { - txnErr := c.withTxn(func() error { - select { - case <-ctx.Done(): - if ctxErr := ctx.Err(); !errors.Is(ctxErr, context.Canceled) { - return ctxErr - } - return nil - case c.Producer.Input() <- c.cascadeMessage(msg): - c.Logger.Log(ctx, slog.LevelInfo, "cascaded", "from", msg.Topic, "to", c.NextTopic) - return nil - } - }) - if txnErr != nil { - c.Logger.Log(ctx, slog.LevelInfo, "Unable to complete cascading transaction", "error", err) - return err - } - } - return nil -} - -// withTxn wraps a function with a transaction that is aborted if an error is returned. -func (c *CascadingConsumer) withTxn(f func() error) (err error) { - if err := c.Producer.BeginTxn(); err != nil { - return fmt.Errorf("unable to begin transaction: %w", err) - } - defer func(err *error) { - if err != nil && *err != nil { - if abortErr := c.Producer.AbortTxn(); abortErr != nil { - c.Logger.Log(nil, slog.LevelInfo, "Unable to abort transaction", "error", abortErr) - } - return - } - if commitErr := c.Producer.CommitTxn(); commitErr != nil { - c.Logger.Log(nil, slog.LevelInfo, "Unable to commit transaction", "error", commitErr) +// CappedExponentialBinaryDelay builds delay functions that use exponential +// binary backoff with a maximum duration. +func CappedExponentialBinaryDelay(cap time.Duration) func(int) time.Duration { + return func(tries int) time.Duration { + b := DelayExponentialBinary(tries) + if b > cap { + return cap } - }(&err) - return f() -} - -// cascadeMessage to the next topic. -func (c *CascadingConsumer) cascadeMessage(msg *sarama.ConsumerMessage) *sarama.ProducerMessage { - pHeaders := make([]sarama.RecordHeader, len(msg.Headers)) - for idx, header := range msg.Headers { - pHeaders[idx] = *header - } - return &sarama.ProducerMessage{ - Key: sarama.ByteEncoder(msg.Key), - Value: sarama.ByteEncoder(msg.Value), - Topic: c.NextTopic, - Headers: c.updateCascadeHeaders(pHeaders), + return b } } -// updateCascadeHeaders calculates not before and failures header values. -// -// Existing not before and failures headers will be dropped in place of the new ones. -func (c *CascadingConsumer) updateCascadeHeaders(headers []sarama.RecordHeader) []sarama.RecordHeader { - failures := 0 - notBefore := time.Now() - - keep := make([]sarama.RecordHeader, 0, len(headers)) - for _, header := range headers { - switch { - case bytes.Equal(header.Key, HeaderNotBefore): - continue // Drop this header, we'll add a new version below. - case bytes.Equal(header.Key, HeaderFailures): - parsed, err := strconv.ParseInt(string(header.Value), 10, 32) - if err != nil { - c.Logger.Log(nil, slog.LevelInfo, "Unable to parse consumption failures count", "error", err) - } else { - failures = int(parsed) - notBefore = notBefore.Add(FailuresToDelay[failures]) - } - continue // Drop this header, we'll add a new version below. - } - keep = append(keep, header) - } - - keep = append(keep, sarama.RecordHeader{ - Key: HeaderNotBefore, - Value: []byte(notBefore.Format(NotBeforeTimeFormat)), - }) - keep = append(keep, sarama.RecordHeader{ - Key: HeaderFailures, - Value: []byte(strconv.Itoa(failures + 1)), - }) - - return keep +// LimitedAsyncProducer restricts the [sarama.AsyncProducer] interface to ensure that its +// recipient isn't able to call Close(), thereby opening the potential for a panic when +// writing to a closed channel. +type LimitedAsyncProducer interface { + AbortTxn() error + BeginTxn() error + CommitTxn() error + Input() chan<- *sarama.ProducerMessage } diff --git a/vendor/github.com/tidepool-org/go-common/asyncevents/sarama.go b/vendor/github.com/tidepool-org/go-common/asyncevents/sarama.go index 83ad151e..0a9a1e31 100644 --- a/vendor/github.com/tidepool-org/go-common/asyncevents/sarama.go +++ b/vendor/github.com/tidepool-org/go-common/asyncevents/sarama.go @@ -11,29 +11,29 @@ import ( "github.com/IBM/sarama" ) -// SaramaEventsConsumer consumes Kafka messages for asynchronous event +// SaramaConsumerGroupManager manages a consumer group for asynchronous Kafka event // handling. -type SaramaEventsConsumer struct { +type SaramaConsumerGroupManager struct { Handler sarama.ConsumerGroupHandler ConsumerGroup sarama.ConsumerGroup Topics []string } -func NewSaramaEventsConsumer(consumerGroup sarama.ConsumerGroup, - handler sarama.ConsumerGroupHandler, topics ...string) *SaramaEventsConsumer { +func NewSaramaConsumerGroupManager(consumerGroup sarama.ConsumerGroup, + handler sarama.ConsumerGroupHandler, topics ...string) *SaramaConsumerGroupManager { - return &SaramaEventsConsumer{ + return &SaramaConsumerGroupManager{ ConsumerGroup: consumerGroup, Handler: handler, Topics: topics, } } -// Run the consumer, to begin consuming Kafka messages. +// Run the manager, to begin consuming Kafka messages. // // Run is stopped by its context being canceled. When its context is canceled, // it returns nil. -func (p *SaramaEventsConsumer) Run(ctx context.Context) (err error) { +func (p *SaramaConsumerGroupManager) Run(ctx context.Context) (err error) { for { err := p.ConsumerGroup.Consume(ctx, p.Topics, p.Handler) if err != nil { @@ -114,9 +114,9 @@ func (h *SaramaConsumerGroupHandler) ConsumeClaim(session sarama.ConsumerGroupSe // Close implements sarama.ConsumerGroupHandler. func (h *SaramaConsumerGroupHandler) Close() error { return nil } -// SaramaMessageConsumer processes Kafka messages. +// SaramaMessageConsumer is responsible for the processing of Kafka messages. type SaramaMessageConsumer interface { - // Consume should process a message. + // Consume processes a message. // // Consume is responsible for marking the message consumed, unless the // context is canceled, in which case the caller should retry, or mark the @@ -126,14 +126,11 @@ type SaramaMessageConsumer interface { var ErrRetriesLimitExceeded = errors.New("retry limit exceeded") -// NTimesRetryingConsumer enhances a SaramaMessageConsumer with a finite -// number of immediate retries. +// NTimesRetryingConsumer is a SaramaMessageConsumer with a finite number of retries. // // The delay between each retry can be controlled via the Delay property. If // no Delay property is specified, a delay based on the Fibonacci sequence is // used. -// -// Logger is intentionally minimal. The slog.Log function is used by default. type NTimesRetryingConsumer struct { Times int Consumer SaramaMessageConsumer @@ -148,6 +145,7 @@ type Logger interface { Log(ctx context.Context, level slog.Level, msg string, args ...any) } +// Consume implements SaramaMessageConsumer. func (c *NTimesRetryingConsumer) Consume(ctx context.Context, session sarama.ConsumerGroupSession, message *sarama.ConsumerMessage) (err error) { diff --git a/vendor/github.com/tidepool-org/go-common/asyncevents/startstopadapter.go b/vendor/github.com/tidepool-org/go-common/asyncevents/startstopadapter.go new file mode 100644 index 00000000..b6e3745b --- /dev/null +++ b/vendor/github.com/tidepool-org/go-common/asyncevents/startstopadapter.go @@ -0,0 +1,124 @@ +package asyncevents + +import ( + "context" + "fmt" + "sync" +) + +// Runner abstracts [SaramaConsumerGroupManager], which is intended to be extended with +// additional capabilities and behaviors. +type Runner interface { + Run(context.Context) error +} + +// BlockingStartStopAdapter for [SaramaConsumerGroupManager] to use Start and Stop methods. +// +// This adapter provides a base for more specific adaptation to adjust behavior to their +// needs. For example [NonBlockingStartStopAdapter] fits well with Uber's fx.Lifecycle, +// while this adapter is more adaptable to platform's event.Runner interface. +type BlockingStartStopAdapter struct { + Runner Runner + + cancelMu sync.Mutex + cancelFunc context.CancelFunc +} + +func NewBlockingStartStopAdapter(runner Runner) *BlockingStartStopAdapter { + return &BlockingStartStopAdapter{ + Runner: runner, + } +} + +func (a *BlockingStartStopAdapter) Start(ctx context.Context) error { + cancelCtx, err := a.init(ctx) + if err != nil { + return err + } + return a.Runner.Run(cancelCtx) +} + +func (a *BlockingStartStopAdapter) init(ctx context.Context) (context.Context, error) { + a.cancelMu.Lock() + defer a.cancelMu.Unlock() + + if a.cancelFunc != nil { + return nil, fmt.Errorf("can't start consumer, it's already running") + } + cancelCtx, cancelFunc := context.WithCancel(ctx) + a.cancelFunc = cancelFunc + return cancelCtx, nil +} + +func (a *BlockingStartStopAdapter) Stop(_ context.Context) error { + a.cancelMu.Lock() + defer a.cancelMu.Unlock() + + if a.cancelFunc == nil { + return fmt.Errorf("can't stop consumer, it's not running") + } + + a.cancelFunc() + a.cancelFunc = nil + + return nil +} + +// NonBlockingStartStopAdapter for [SaramaConsumerGroupManager] for non-blocking Start and +// Stop methods. +// +// To facilitate error reporting during non-blocking operation, a callback can be provided, +// which if defined, will be called with errors that cause a [SaramaConsumerGroupManager]'s +// Run method to return. In addition, when the callback is defined, panics from within Run +// are recovered, converted to errors, and passed to the callback before being discarded. +type NonBlockingStartStopAdapter struct { + *BlockingStartStopAdapter + + onError func(error) +} + +func NewNonBlockingStartStopAdapter(consumer Runner, onError func(error)) *NonBlockingStartStopAdapter { + blocking := NewBlockingStartStopAdapter(consumer) + return &NonBlockingStartStopAdapter{ + BlockingStartStopAdapter: blocking, + onError: onError, + } +} + +func (a *NonBlockingStartStopAdapter) Start(ctx context.Context) error { + go a.start(ctx) + return nil +} + +func (a *NonBlockingStartStopAdapter) start(ctx context.Context) { + defer a.maybeRecover() + + if err := a.BlockingStartStopAdapter.Start(ctx); err != nil { + if a.onError != nil { + a.onError(err) + } + } +} + +// maybeRecover uses a callback, if defined, to process recovered panics. +// +// If the callback isn't defined, the panic will be re-raised. +func (a *NonBlockingStartStopAdapter) maybeRecover() { + if r := recover(); r != nil { + if a.onError != nil { + a.onError(a.wrapPanic(r)) + } else { + panic(r) + } + } +} + +// wrapPanic converts a non-error value to an error for passing to an on-error callback. +// +// Existing error values are wrapped, to allow later unwrapping. +func (a *NonBlockingStartStopAdapter) wrapPanic(r any) error { + if err, ok := r.(error); ok { + return fmt.Errorf("consumer panicked: %w", err) + } + return fmt.Errorf("consumer panicked: %s", r) +} diff --git a/vendor/modules.txt b/vendor/modules.txt index e4e229a1..73a19b4a 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -207,7 +207,7 @@ github.com/tidepool-org/clinic/client # github.com/tidepool-org/clinic/redox_models v0.0.0-20250521172904-b61821ac9973 ## explicit; go 1.22 github.com/tidepool-org/clinic/redox_models -# github.com/tidepool-org/go-common v0.12.3-0.20250613120630-656deb326ad3 +# github.com/tidepool-org/go-common v0.12.3-0.20250625233514-f3a3357c3d91 ## explicit; go 1.24.1 github.com/tidepool-org/go-common/asyncevents github.com/tidepool-org/go-common/clients diff --git a/worker/bootstrap.go b/worker/bootstrap.go index 774b551d..e66b0848 100644 --- a/worker/bootstrap.go +++ b/worker/bootstrap.go @@ -1,10 +1,12 @@ package worker import ( + "net/http" + "github.com/tidepool-org/clinic-worker/merge" "github.com/tidepool-org/clinic-worker/redox" - "github.com/tidepool-org/go-common/asyncevents" - "net/http" + + "go.uber.org/fx" "github.com/tidepool-org/clinic-worker/cdc" "github.com/tidepool-org/clinic-worker/clinicians" @@ -16,7 +18,6 @@ import ( "github.com/tidepool-org/clinic-worker/patientsummary" "github.com/tidepool-org/clinic-worker/users" "github.com/tidepool-org/go-common/events" - "go.uber.org/fx" ) var dependencies = fx.Provide( @@ -47,6 +48,7 @@ var Modules = []fx.Option{ redox.Module, users.Module, marketo.Module, + cdc.Module, } func New() *fx.App { @@ -60,8 +62,7 @@ func New() *fx.App { type Components struct { fx.In - Consumers []events.EventConsumer `group:"consumers"` - Runners []asyncevents.SaramaEventsRunner `group:"runners"` + Consumers []events.EventConsumer `group:"consumers"` HealthCheckServer *http.Server Lifecycle fx.Lifecycle Shutdowner fx.Shutdowner @@ -71,7 +72,4 @@ func startConsumers(components Components) { for _, consumer := range components.Consumers { cdc.AttachConsumerGroupHooks(consumer, components.Lifecycle, components.Shutdowner) } - for _, runner := range components.Runners { - cdc.AttachSaramaRunnerHooks(runner, components.Lifecycle, components.Shutdowner) - } }