From 00501d4dae64d9c9f8fac5539c72164f9a1f5403 Mon Sep 17 00:00:00 2001 From: Bastien BEUCHAT Date: Fri, 3 Jul 2026 12:07:04 +0200 Subject: [PATCH] [CLOUDTRUST-8743] Health check : Add kafka consumer + remove audit --- healthcheck/audit.go | 63 ---------------- healthcheck/audit_test.go | 44 ------------ healthcheck/health.go | 9 ++- healthcheck/kafkaconsumer.go | 47 ++++++++++++ healthcheck/kafkaconsumer_test.go | 92 ++++++++++++++++++++++++ healthcheck/mock/eventsreportermodule.go | 54 -------------- healthcheck/mock/kafka.go | 54 ++++++++++++++ healthcheck/mock_test.go | 2 +- 8 files changed, 198 insertions(+), 167 deletions(-) delete mode 100644 healthcheck/audit.go delete mode 100644 healthcheck/audit_test.go create mode 100644 healthcheck/kafkaconsumer.go create mode 100644 healthcheck/kafkaconsumer_test.go delete mode 100644 healthcheck/mock/eventsreportermodule.go create mode 100644 healthcheck/mock/kafka.go diff --git a/healthcheck/audit.go b/healthcheck/audit.go deleted file mode 100644 index 45bd3ff..0000000 --- a/healthcheck/audit.go +++ /dev/null @@ -1,63 +0,0 @@ -package healthcheck - -import ( - "context" - "time" - - "github.com/cloudtrust/common-service/v2/events" - log "github.com/cloudtrust/common-service/v2/log" -) - -type auditEventsReporterChecker struct { - alias string - reporter events.AuditEventsReporterModule - timeout time.Duration - response HealthStatus - logger log.Logger - failureCounter int -} - -func newAuditEventsReporterChecker(alias string, reporter events.AuditEventsReporterModule, timeout time.Duration, cacheDuration time.Duration, logger log.Logger, timeProvider TimeProvider) BasicChecker { - healthStatusType := "auditEventreporter" - response := HealthStatus{Name: &alias, Type: &healthStatusType, CacheDuration: cacheDuration, TimeProvider: timeProvider} - response.connection("init") - response.stateUp() - return &auditEventsReporterChecker{ - alias: alias, - reporter: reporter, - timeout: timeout, - response: response, - logger: logger, - failureCounter: 0, - } -} - -func (a *auditEventsReporterChecker) CheckStatus() HealthStatus { - if a.response.hasExpired() { - go a.updateStatus() - } - - return a.response -} - -func (a *auditEventsReporterChecker) updateStatus() { - finished := make(chan bool) - go func() { - event := events.NewEvent("healthcheck", "", "master", "health_checker", "health_checker", "master", map[string]string{}) - a.reporter.ReportEvent(context.Background(), event) - finished <- true - }() - - select { - case <-finished: - a.response.connection("established") - a.response.stateUp() - a.failureCounter = 0 - case <-time.After(a.timeout): - a.response.stateDown("Events reporter timeout") - a.failureCounter++ - a.logger.Error(context.Background(), "msg", "Audit Events Reporter timeout to produce events", "timeout", a.timeout, "failureCounter", a.failureCounter) - } - - a.response.touch() -} diff --git a/healthcheck/audit_test.go b/healthcheck/audit_test.go deleted file mode 100644 index 81ac48b..0000000 --- a/healthcheck/audit_test.go +++ /dev/null @@ -1,44 +0,0 @@ -package healthcheck - -import ( - "testing" - "time" - - "github.com/cloudtrust/common-service/v2/healthcheck/mock" - log "github.com/cloudtrust/common-service/v2/log" - "github.com/stretchr/testify/assert" - "go.uber.org/mock/gomock" -) - -func TestAuditEventsReporterChecker(t *testing.T) { - var mockCtrl = gomock.NewController(t) - defer mockCtrl.Finish() - - var mockAuditEventReporter = mock.NewAuditEventsReporterModule(mockCtrl) - var mockTime = mock.NewTimeProvider(mockCtrl) - mockTime.EXPECT().Now().Return(testTime).AnyTimes() - - t.Run("Success ", func(t *testing.T) { - var auditEventReporterChecker = newAuditEventsReporterChecker("alias", mockAuditEventReporter, 10*time.Second, 10*time.Second, log.NewNopLogger(), mockTime) - internalChecker := auditEventReporterChecker.(*auditEventsReporterChecker) - mockAuditEventReporter.EXPECT().ReportEvent(gomock.Any(), gomock.Any()).Times(1) - - internalChecker.updateStatus() - var res = internalChecker.response - assert.NotNil(t, res.Connection) - assert.Equal(t, "established", *res.Connection) - }) - - t.Run("Failure ", func(t *testing.T) { - var auditEventReporterChecker = newAuditEventsReporterChecker("alias", mockAuditEventReporter, 1*time.Second, 10*time.Second, log.NewNopLogger(), mockTime) - internalChecker := auditEventReporterChecker.(*auditEventsReporterChecker) - mockAuditEventReporter.EXPECT().ReportEvent(gomock.Any(), gomock.Any()).Do(func(arg0 any, arg1 any) { - time.Sleep(2 * time.Second) - }) - - internalChecker.updateStatus() - res := internalChecker.response - assert.NotNil(t, res.Message) - assert.Equal(t, "Events reporter timeout", *res.Message) - }) -} diff --git a/healthcheck/health.go b/healthcheck/health.go index f3d8cd8..b36015e 100644 --- a/healthcheck/health.go +++ b/healthcheck/health.go @@ -5,7 +5,6 @@ import ( "net/http" "time" - "github.com/cloudtrust/common-service/v2/events" commonhttp "github.com/cloudtrust/common-service/v2/http" log "github.com/cloudtrust/common-service/v2/log" "github.com/go-kit/kit/ratelimit" @@ -18,7 +17,7 @@ type HealthChecker interface { AddHTTPEndpoint(name string, targetURL string, timeoutDuration time.Duration, expectedStatus int, cacheDuration time.Duration) AddHTTPEndpoints(endpoints map[string]string, timeoutDuration time.Duration, expectedStatus int, cacheDuration time.Duration) AddDatabase(name string, db HealthDatabase, cacheDuration time.Duration) - AddAuditEventsReporterModule(name string, reporter events.AuditEventsReporterModule, timeout time.Duration, cacheDuration time.Duration) + AddKafkaConsumer(name string, consumer KafkaConsumer, cacheDuration time.Duration) MakeHandler(rateLimit ratelimit.Allower) http.HandlerFunc } @@ -134,9 +133,9 @@ func (hc *healthchecker) AddDatabase(name string, db HealthDatabase, cacheDurati hc.AddHealthChecker(name, newDatabaseChecker(name, db, cacheDuration, RealTimeProvider{})) } -func (hc *healthchecker) AddAuditEventsReporterModule(name string, reporter events.AuditEventsReporterModule, timeout time.Duration, cacheDuration time.Duration) { - hc.logger.Info(context.Background(), "msg", "Adding audit event reporter module", "processor", name) - hc.AddHealthChecker(name, newAuditEventsReporterChecker(name, reporter, timeout, cacheDuration, hc.logger, RealTimeProvider{})) +func (hc *healthchecker) AddKafkaConsumer(name string, consumer KafkaConsumer, cacheDuration time.Duration) { + hc.logger.Info(context.Background(), "msg", "Adding Kafka consumer", "processor", name) + hc.AddHealthChecker(name, newKafkaConsumerChecker(name, consumer, cacheDuration, RealTimeProvider{})) } // MakeHandler makes a HTTP handler that returns health check information diff --git a/healthcheck/kafkaconsumer.go b/healthcheck/kafkaconsumer.go new file mode 100644 index 0000000..124cae2 --- /dev/null +++ b/healthcheck/kafkaconsumer.go @@ -0,0 +1,47 @@ +package healthcheck + +import ( + "time" +) + +type KafkaConsumer interface { + IsLive() bool +} + +type kafkaConsumerChecker struct { + alias string + consumer KafkaConsumer + response HealthStatus + failureCounter int +} + +func newKafkaConsumerChecker(alias string, consumer KafkaConsumer, cacheDuration time.Duration, timeProvider TimeProvider) BasicChecker { + healthStatusType := "kafkaConsumer" + response := HealthStatus{Name: &alias, Type: &healthStatusType, CacheDuration: cacheDuration, TimeProvider: timeProvider} + response.connection("init") + response.stateUp() + return &kafkaConsumerChecker{ + alias: alias, + consumer: consumer, + response: response, + failureCounter: 0, + } +} + +func (k *kafkaConsumerChecker) CheckStatus() HealthStatus { + if !k.response.hasExpired() { + return k.response + } + + if k.consumer.IsLive() { + k.response.connection("consuming") + k.response.stateUp() + k.failureCounter = 0 + } else { + k.response.stateDown("Kafka consumer is not live") + k.failureCounter++ + } + + k.response.touch() + return k.response +} diff --git a/healthcheck/kafkaconsumer_test.go b/healthcheck/kafkaconsumer_test.go new file mode 100644 index 0000000..af28897 --- /dev/null +++ b/healthcheck/kafkaconsumer_test.go @@ -0,0 +1,92 @@ +package healthcheck + +import ( + "testing" + "time" + + "github.com/cloudtrust/common-service/v2/healthcheck/mock" + "github.com/stretchr/testify/assert" + "go.uber.org/mock/gomock" +) + +func TestKafkaConsumerHealthCheckLiveAndCached(t *testing.T) { + mockCtrl := gomock.NewController(t) + defer mockCtrl.Finish() + + mockTime := mock.NewTimeProvider(mockCtrl) + mockTime.EXPECT().Now().Return(testTime).AnyTimes() + + consumer := mock.NewKafkaConsumer(mockCtrl) + consumer.EXPECT().IsLive().Return(true).Times(1) + checker := newKafkaConsumerChecker("alias", consumer, 10*time.Second, mockTime) + + status := checker.CheckStatus() + assert.Equal(t, "UP", *status.State) + assert.Equal(t, "consuming", *status.Connection) + assert.Nil(t, status.Message) + + status = checker.CheckStatus() + assert.Equal(t, "UP", *status.State) + assert.Equal(t, "consuming", *status.Connection) + + internal := checker.(*kafkaConsumerChecker) + assert.Equal(t, 0, internal.failureCounter) +} + +func TestKafkaConsumerHealthCheckDownAndCached(t *testing.T) { + mockCtrl := gomock.NewController(t) + defer mockCtrl.Finish() + + mockTime := mock.NewTimeProvider(mockCtrl) + mockTime.EXPECT().Now().Return(testTime).AnyTimes() + + consumer := mock.NewKafkaConsumer(mockCtrl) + consumer.EXPECT().IsLive().Return(false).Times(1) + checker := newKafkaConsumerChecker("alias", consumer, 10*time.Second, mockTime) + + status := checker.CheckStatus() + assert.Equal(t, "DOWN", *status.State) + assert.Nil(t, status.Connection) + assert.NotNil(t, status.Message) + assert.Equal(t, "Kafka consumer is not live", *status.Message) + + status = checker.CheckStatus() + assert.Equal(t, "DOWN", *status.State) + assert.Equal(t, "Kafka consumer is not live", *status.Message) + + internal := checker.(*kafkaConsumerChecker) + assert.Equal(t, 1, internal.failureCounter) +} + +func TestKafkaConsumerHealthCheckFailureCounterResetOnRecovery(t *testing.T) { + mockCtrl := gomock.NewController(t) + defer mockCtrl.Finish() + + mockTime := mock.NewTimeProvider(mockCtrl) + gomock.InOrder( + mockTime.EXPECT().Now().Return(testTime), + mockTime.EXPECT().Now().Return(testTime), + mockTime.EXPECT().Now().Return(testTime.Add(10*time.Second).Add(time.Millisecond)), + mockTime.EXPECT().Now().Return(testTime.Add(10*time.Second).Add(time.Millisecond)), + ) + + consumer := mock.NewKafkaConsumer(mockCtrl) + gomock.InOrder( + consumer.EXPECT().IsLive().Return(false), + consumer.EXPECT().IsLive().Return(true), + ) + checker := newKafkaConsumerChecker("alias", consumer, 10*time.Second, mockTime) + + status := checker.CheckStatus() + assert.Equal(t, "DOWN", *status.State) + assert.Equal(t, "Kafka consumer is not live", *status.Message) + + internal := checker.(*kafkaConsumerChecker) + assert.Equal(t, 1, internal.failureCounter) + + status = checker.CheckStatus() + assert.Equal(t, "UP", *status.State) + assert.Equal(t, "consuming", *status.Connection) + assert.Nil(t, status.Message) + assert.Equal(t, 0, internal.failureCounter) +} diff --git a/healthcheck/mock/eventsreportermodule.go b/healthcheck/mock/eventsreportermodule.go deleted file mode 100644 index 4c1e2b3..0000000 --- a/healthcheck/mock/eventsreportermodule.go +++ /dev/null @@ -1,54 +0,0 @@ -// Code generated by MockGen. DO NOT EDIT. -// Source: github.com/cloudtrust/common-service/v2/events (interfaces: AuditEventsReporterModule) -// -// Generated by this command: -// -// mockgen --build_flags=--mod=mod -destination=./mock/eventsreportermodule.go -package=mock -mock_names=AuditEventsReporterModule=AuditEventsReporterModule github.com/cloudtrust/common-service/v2/events AuditEventsReporterModule -// - -// Package mock is a generated GoMock package. -package mock - -import ( - context "context" - reflect "reflect" - - events "github.com/cloudtrust/common-service/v2/events" - gomock "go.uber.org/mock/gomock" -) - -// AuditEventsReporterModule is a mock of AuditEventsReporterModule interface. -type AuditEventsReporterModule struct { - ctrl *gomock.Controller - recorder *AuditEventsReporterModuleMockRecorder - isgomock struct{} -} - -// AuditEventsReporterModuleMockRecorder is the mock recorder for AuditEventsReporterModule. -type AuditEventsReporterModuleMockRecorder struct { - mock *AuditEventsReporterModule -} - -// NewAuditEventsReporterModule creates a new mock instance. -func NewAuditEventsReporterModule(ctrl *gomock.Controller) *AuditEventsReporterModule { - mock := &AuditEventsReporterModule{ctrl: ctrl} - mock.recorder = &AuditEventsReporterModuleMockRecorder{mock} - return mock -} - -// EXPECT returns an object that allows the caller to indicate expected use. -func (m *AuditEventsReporterModule) EXPECT() *AuditEventsReporterModuleMockRecorder { - return m.recorder -} - -// ReportEvent mocks base method. -func (m *AuditEventsReporterModule) ReportEvent(ctx context.Context, event events.Event) { - m.ctrl.T.Helper() - m.ctrl.Call(m, "ReportEvent", ctx, event) -} - -// ReportEvent indicates an expected call of ReportEvent. -func (mr *AuditEventsReporterModuleMockRecorder) ReportEvent(ctx, event any) *gomock.Call { - mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ReportEvent", reflect.TypeOf((*AuditEventsReporterModule)(nil).ReportEvent), ctx, event) -} diff --git a/healthcheck/mock/kafka.go b/healthcheck/mock/kafka.go new file mode 100644 index 0000000..1398bf5 --- /dev/null +++ b/healthcheck/mock/kafka.go @@ -0,0 +1,54 @@ +// Code generated by MockGen. DO NOT EDIT. +// Source: github.com/cloudtrust/common-service/v2/healthcheck (interfaces: KafkaConsumer) +// +// Generated by this command: +// +// mockgen --build_flags=--mod=mod -destination=./mock/kafka.go -package=mock -mock_names=KafkaConsumer=KafkaConsumer github.com/cloudtrust/common-service/v2/healthcheck KafkaConsumer +// + +// Package mock is a generated GoMock package. +package mock + +import ( + reflect "reflect" + + gomock "go.uber.org/mock/gomock" +) + +// KafkaConsumer is a mock of KafkaConsumer interface. +type KafkaConsumer struct { + ctrl *gomock.Controller + recorder *KafkaConsumerMockRecorder + isgomock struct{} +} + +// KafkaConsumerMockRecorder is the mock recorder for KafkaConsumer. +type KafkaConsumerMockRecorder struct { + mock *KafkaConsumer +} + +// NewKafkaConsumer creates a new mock instance. +func NewKafkaConsumer(ctrl *gomock.Controller) *KafkaConsumer { + mock := &KafkaConsumer{ctrl: ctrl} + mock.recorder = &KafkaConsumerMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use. +func (m *KafkaConsumer) EXPECT() *KafkaConsumerMockRecorder { + return m.recorder +} + +// IsLive mocks base method. +func (m *KafkaConsumer) IsLive() bool { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "IsLive") + ret0, _ := ret[0].(bool) + return ret0 +} + +// IsLive indicates an expected call of IsLive. +func (mr *KafkaConsumerMockRecorder) IsLive() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "IsLive", reflect.TypeOf((*KafkaConsumer)(nil).IsLive)) +} diff --git a/healthcheck/mock_test.go b/healthcheck/mock_test.go index 9857a21..a8a2d54 100644 --- a/healthcheck/mock_test.go +++ b/healthcheck/mock_test.go @@ -1,5 +1,5 @@ package healthcheck //go:generate mockgen --build_flags=--mod=mod -destination=./mock/healthcheck.go -package=mock -mock_names=HealthDatabase=HealthDatabase github.com/cloudtrust/common-service/v2/healthcheck HealthDatabase -//go:generate mockgen --build_flags=--mod=mod -destination=./mock/eventsreportermodule.go -package=mock -mock_names=AuditEventsReporterModule=AuditEventsReporterModule github.com/cloudtrust/common-service/v2/events AuditEventsReporterModule //go:generate mockgen --build_flags=--mod=mod -destination=./mock/timeprovider.go -package=mock -mock_names=TimeProvider=TimeProvider github.com/cloudtrust/common-service/v2/healthcheck TimeProvider +//go:generate mockgen --build_flags=--mod=mod -destination=./mock/kafka.go -package=mock -mock_names=KafkaConsumer=KafkaConsumer github.com/cloudtrust/common-service/v2/healthcheck KafkaConsumer