Skip to content
Open
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
100 changes: 77 additions & 23 deletions cdc/init.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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)
}),
),
)
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
42 changes: 23 additions & 19 deletions redox/consumer_message.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand All @@ -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
}
Expand All @@ -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 {
Expand Down
40 changes: 22 additions & 18 deletions redox/consumer_scheduled.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand All @@ -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
}
Expand All @@ -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
Expand Down
6 changes: 4 additions & 2 deletions redox/module.go
Original file line number Diff line number Diff line change
@@ -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(
Expand Down
Loading