diff --git a/.gitlab/.gitlab-ci.yml b/.gitlab/.gitlab-ci.yml index b2e84e7a8..3a716508b 100644 --- a/.gitlab/.gitlab-ci.yml +++ b/.gitlab/.gitlab-ci.yml @@ -139,7 +139,7 @@ build amd container fb: before_script: - echo "${DOCKER_DEPLOY_TOKEN}" | docker login -u ${DOCKER_DEPLOY_USER} --password-stdin ${CI_REGISTRY} script: - - tag=${CI_REGISTRY_IMAGE}/featurebase:linux-amd64-${CI_COMMIT_REF_SLUG} + - tag=${CI_REGISTRY_IMAGE}/featurebase:linux-amd64-${CI_COMMIT_REF_NAME} - docker build --build-arg GO_VERSION=$GOVERSION --build-arg ARCH=amd64 -t $tag -f .gitlab/Dockerfile . - docker push $tag - echo Created docker featurebase image with tag "$tag" @@ -578,7 +578,7 @@ build arm container fb: before_script: - echo "${DOCKER_DEPLOY_TOKEN}" | docker login -u ${DOCKER_DEPLOY_USER} --password-stdin ${CI_REGISTRY} script: - - tag=${CI_REGISTRY_IMAGE}/featurebase:linux-arm64-${CI_COMMIT_REF_SLUG} + - tag=${CI_REGISTRY_IMAGE}/featurebase:linux-arm64-${CI_COMMIT_REF_NAME} - docker build --build-arg GO_VERSION=$GOVERSION --build-arg ARCH=arm64 -t $tag -f .gitlab/Dockerfile . - docker push $tag - echo Created docker featurebase image with tag "$tag" diff --git a/go.mod b/go.mod index d5a36c5de..6660519b2 100644 --- a/go.mod +++ b/go.mod @@ -73,6 +73,7 @@ require ( ) require ( + github.com/google/uuid v1.3.0 github.com/jaffee/commandeer v0.5.0 github.com/linkedin/goavro/v2 v2.11.1 google.golang.org/grpc v1.46.0 @@ -97,7 +98,6 @@ require ( github.com/golang/snappy v0.0.4 // indirect github.com/google/btree v1.0.1 // indirect github.com/google/pprof v0.0.0-20211214055906-6f57359322fd // indirect - github.com/google/uuid v1.3.0 // indirect github.com/gopherjs/gopherjs v0.0.0-20200217142428-fce0ec30dd00 // indirect github.com/gorilla/websocket v1.5.0 // indirect github.com/grpc-ecosystem/go-grpc-middleware v1.3.0 // indirect diff --git a/idk/Makefile b/idk/Makefile index bcb429902..100aa8acf 100644 --- a/idk/Makefile +++ b/idk/Makefile @@ -1,6 +1,7 @@ GO ?= go GOOS ?= $(shell $(GO) env GOOS) GOARCH ?= $(shell $(GO) env GOARCH) +GOPATH ?= $(shell $(GO) env GOPATH) GO_VERSION=1.17.8 GO_BUILD_FLAGS= ODBC_ENABLED= @@ -309,3 +310,21 @@ docker-push-ecr: docker-image aws-login aws-login: aws sso login --profile $(PROFILE) + + +# Build mock implementations of AWS service API interfaces for unit testing. +# +# These are run-once/rarely targets as they produce code artifacts only used +# during `go test`. An automated CI/CD process does NOT need to invoke every time. +MOCKS = S3 Kinesis SQS +install-mock-generator: + @which $(GOPATH)/bin/mockery || (echo "Installing missing dependency 'mockery' to generate mock implementations for unit testing." && \ + $(GO) install 'github.com/vektra/mockery/v2@latest') + +update-mocks: $(addprefix update-mocks-, $(MOCKS)) + +update-mocks-%: install-mock-generator + $(eval AWS_SERVICE := $(shell echo $* | tr '[:upper:]' '[:lower:]')) + $(eval AWS_SDK_VERSION := $(shell grep 'github.com/aws/aws-sdk-go' ../go.mod | cut -d ' ' -f 2)) + echo Generating mock for AWS service $(AWS_SERVICE) and SDK version $(AWS_SDK_VERSION) && \ + $(GOPATH)/bin/mockery --name $*API --output idktest/mocks --filename $(AWS_SERVICE).go --dir $(GOPATH)/pkg/mod/github.com/aws/aws-sdk-go@$(AWS_SDK_VERSION)/service/$(AWS_SERVICE)/$(AWS_SERVICE)iface diff --git a/idk/cmd/molecula-consumer-kinesis/main.go b/idk/cmd/molecula-consumer-kinesis/main.go index f445137f0..3eb1f8af3 100644 --- a/idk/cmd/molecula-consumer-kinesis/main.go +++ b/idk/cmd/molecula-consumer-kinesis/main.go @@ -9,24 +9,42 @@ import ( "github.com/molecula/featurebase/v3/logger" ) +func logFailure(errorType kinesis.ErrorType, m *kinesis.Main, v interface{}) { + log := m.Log() + + if log == nil { + log = logger.NewStandardLogger(os.Stderr) + } + + if errorType == kinesis.RecoverableErrorType { + log.Errorf("Error running command: %+v", v) + } else { + log.Panicf("Panic running command: %+v", v) + } +} + func main() { m := kinesis.NewMain() if err := pflag.LoadEnv(m, "CONSUMER_", nil); err != nil { log.Fatal(err) } m.Rename() + + // Capture any panic and log it before dying. + defer func() { + if r := recover(); r != nil { + logFailure(kinesis.PanicErrorType, m, r) + os.Exit(1) + } + }() + if m.DryRun { log.Printf("%+v\n", m) return } + if err := m.Run(); err != nil { - log := m.Log() - if log == nil { - // if we fail before a logger was instantiated - logger.NewStandardLogger(os.Stderr).Errorf("Error running command: %v", err) - os.Exit(1) - } - log.Errorf("Error running command: %v", err) + logFailure(kinesis.RecoverableErrorType, m, err) os.Exit(1) } } diff --git a/idk/idktest/mocks/kinesis.go b/idk/idktest/mocks/kinesis.go index 4cbbec253..063cf51ed 100644 --- a/idk/idktest/mocks/kinesis.go +++ b/idk/idktest/mocks/kinesis.go @@ -1,9 +1,7 @@ -// Code generated by mockery v1.0.0. DO NOT EDIT. +// Code generated by mockery v2.14.0. DO NOT EDIT. package mocks -import "github.com/aws/aws-sdk-go/service/kinesis/kinesisiface" - import ( context "context" @@ -2456,4 +2454,17 @@ func (_m *KinesisAPI) WaitUntilStreamNotExistsWithContext(_a0 context.Context, _ return r0 } -var _ kinesisiface.KinesisAPI = (*KinesisAPI)(nil) +type mockConstructorTestingTNewKinesisAPI interface { + mock.TestingT + Cleanup(func()) +} + +// NewKinesisAPI creates a new instance of KinesisAPI. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. +func NewKinesisAPI(t mockConstructorTestingTNewKinesisAPI) *KinesisAPI { + mock := &KinesisAPI{} + mock.Mock.Test(t) + + t.Cleanup(func() { mock.AssertExpectations(t) }) + + return mock +} diff --git a/idk/idktest/mocks/s3.go b/idk/idktest/mocks/s3.go index 5111dba99..883970e1b 100644 --- a/idk/idktest/mocks/s3.go +++ b/idk/idktest/mocks/s3.go @@ -1,9 +1,7 @@ -// Code generated by mockery v1.0.0. DO NOT EDIT. +// Code generated by mockery v2.14.0. DO NOT EDIT. package mocks -import "github.com/aws/aws-sdk-go/service/s3/s3iface" - import ( context "context" @@ -7821,4 +7819,17 @@ func (_m *S3API) WriteGetObjectResponseWithContext(_a0 context.Context, _a1 *s3. return r0, r1 } -var _ s3iface.S3API = (*S3API)(nil) +type mockConstructorTestingTNewS3API interface { + mock.TestingT + Cleanup(func()) +} + +// NewS3API creates a new instance of S3API. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. +func NewS3API(t mockConstructorTestingTNewS3API) *S3API { + mock := &S3API{} + mock.Mock.Test(t) + + t.Cleanup(func() { mock.AssertExpectations(t) }) + + return mock +} diff --git a/idk/idktest/mocks/sqs.go b/idk/idktest/mocks/sqs.go new file mode 100644 index 000000000..33a8b4888 --- /dev/null +++ b/idk/idktest/mocks/sqs.go @@ -0,0 +1,1662 @@ +// Code generated by mockery v2.14.0. DO NOT EDIT. + +package mocks + +import ( + context "context" + + request "github.com/aws/aws-sdk-go/aws/request" + mock "github.com/stretchr/testify/mock" + + sqs "github.com/aws/aws-sdk-go/service/sqs" +) + +// SQSAPI is an autogenerated mock type for the SQSAPI type +type SQSAPI struct { + mock.Mock +} + +// AddPermission provides a mock function with given fields: _a0 +func (_m *SQSAPI) AddPermission(_a0 *sqs.AddPermissionInput) (*sqs.AddPermissionOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.AddPermissionOutput + if rf, ok := ret.Get(0).(func(*sqs.AddPermissionInput) *sqs.AddPermissionOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.AddPermissionOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.AddPermissionInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// AddPermissionRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) AddPermissionRequest(_a0 *sqs.AddPermissionInput) (*request.Request, *sqs.AddPermissionOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.AddPermissionInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.AddPermissionOutput + if rf, ok := ret.Get(1).(func(*sqs.AddPermissionInput) *sqs.AddPermissionOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.AddPermissionOutput) + } + } + + return r0, r1 +} + +// AddPermissionWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) AddPermissionWithContext(_a0 context.Context, _a1 *sqs.AddPermissionInput, _a2 ...request.Option) (*sqs.AddPermissionOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.AddPermissionOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.AddPermissionInput, ...request.Option) *sqs.AddPermissionOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.AddPermissionOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.AddPermissionInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ChangeMessageVisibility provides a mock function with given fields: _a0 +func (_m *SQSAPI) ChangeMessageVisibility(_a0 *sqs.ChangeMessageVisibilityInput) (*sqs.ChangeMessageVisibilityOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.ChangeMessageVisibilityOutput + if rf, ok := ret.Get(0).(func(*sqs.ChangeMessageVisibilityInput) *sqs.ChangeMessageVisibilityOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ChangeMessageVisibilityOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.ChangeMessageVisibilityInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ChangeMessageVisibilityBatch provides a mock function with given fields: _a0 +func (_m *SQSAPI) ChangeMessageVisibilityBatch(_a0 *sqs.ChangeMessageVisibilityBatchInput) (*sqs.ChangeMessageVisibilityBatchOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.ChangeMessageVisibilityBatchOutput + if rf, ok := ret.Get(0).(func(*sqs.ChangeMessageVisibilityBatchInput) *sqs.ChangeMessageVisibilityBatchOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ChangeMessageVisibilityBatchOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.ChangeMessageVisibilityBatchInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ChangeMessageVisibilityBatchRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) ChangeMessageVisibilityBatchRequest(_a0 *sqs.ChangeMessageVisibilityBatchInput) (*request.Request, *sqs.ChangeMessageVisibilityBatchOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.ChangeMessageVisibilityBatchInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.ChangeMessageVisibilityBatchOutput + if rf, ok := ret.Get(1).(func(*sqs.ChangeMessageVisibilityBatchInput) *sqs.ChangeMessageVisibilityBatchOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.ChangeMessageVisibilityBatchOutput) + } + } + + return r0, r1 +} + +// ChangeMessageVisibilityBatchWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) ChangeMessageVisibilityBatchWithContext(_a0 context.Context, _a1 *sqs.ChangeMessageVisibilityBatchInput, _a2 ...request.Option) (*sqs.ChangeMessageVisibilityBatchOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.ChangeMessageVisibilityBatchOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.ChangeMessageVisibilityBatchInput, ...request.Option) *sqs.ChangeMessageVisibilityBatchOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ChangeMessageVisibilityBatchOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.ChangeMessageVisibilityBatchInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ChangeMessageVisibilityRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) ChangeMessageVisibilityRequest(_a0 *sqs.ChangeMessageVisibilityInput) (*request.Request, *sqs.ChangeMessageVisibilityOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.ChangeMessageVisibilityInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.ChangeMessageVisibilityOutput + if rf, ok := ret.Get(1).(func(*sqs.ChangeMessageVisibilityInput) *sqs.ChangeMessageVisibilityOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.ChangeMessageVisibilityOutput) + } + } + + return r0, r1 +} + +// ChangeMessageVisibilityWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) ChangeMessageVisibilityWithContext(_a0 context.Context, _a1 *sqs.ChangeMessageVisibilityInput, _a2 ...request.Option) (*sqs.ChangeMessageVisibilityOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.ChangeMessageVisibilityOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.ChangeMessageVisibilityInput, ...request.Option) *sqs.ChangeMessageVisibilityOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ChangeMessageVisibilityOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.ChangeMessageVisibilityInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// CreateQueue provides a mock function with given fields: _a0 +func (_m *SQSAPI) CreateQueue(_a0 *sqs.CreateQueueInput) (*sqs.CreateQueueOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.CreateQueueOutput + if rf, ok := ret.Get(0).(func(*sqs.CreateQueueInput) *sqs.CreateQueueOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.CreateQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.CreateQueueInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// CreateQueueRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) CreateQueueRequest(_a0 *sqs.CreateQueueInput) (*request.Request, *sqs.CreateQueueOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.CreateQueueInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.CreateQueueOutput + if rf, ok := ret.Get(1).(func(*sqs.CreateQueueInput) *sqs.CreateQueueOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.CreateQueueOutput) + } + } + + return r0, r1 +} + +// CreateQueueWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) CreateQueueWithContext(_a0 context.Context, _a1 *sqs.CreateQueueInput, _a2 ...request.Option) (*sqs.CreateQueueOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.CreateQueueOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.CreateQueueInput, ...request.Option) *sqs.CreateQueueOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.CreateQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.CreateQueueInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// DeleteMessage provides a mock function with given fields: _a0 +func (_m *SQSAPI) DeleteMessage(_a0 *sqs.DeleteMessageInput) (*sqs.DeleteMessageOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.DeleteMessageOutput + if rf, ok := ret.Get(0).(func(*sqs.DeleteMessageInput) *sqs.DeleteMessageOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.DeleteMessageOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.DeleteMessageInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// DeleteMessageBatch provides a mock function with given fields: _a0 +func (_m *SQSAPI) DeleteMessageBatch(_a0 *sqs.DeleteMessageBatchInput) (*sqs.DeleteMessageBatchOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.DeleteMessageBatchOutput + if rf, ok := ret.Get(0).(func(*sqs.DeleteMessageBatchInput) *sqs.DeleteMessageBatchOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.DeleteMessageBatchOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.DeleteMessageBatchInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// DeleteMessageBatchRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) DeleteMessageBatchRequest(_a0 *sqs.DeleteMessageBatchInput) (*request.Request, *sqs.DeleteMessageBatchOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.DeleteMessageBatchInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.DeleteMessageBatchOutput + if rf, ok := ret.Get(1).(func(*sqs.DeleteMessageBatchInput) *sqs.DeleteMessageBatchOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.DeleteMessageBatchOutput) + } + } + + return r0, r1 +} + +// DeleteMessageBatchWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) DeleteMessageBatchWithContext(_a0 context.Context, _a1 *sqs.DeleteMessageBatchInput, _a2 ...request.Option) (*sqs.DeleteMessageBatchOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.DeleteMessageBatchOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.DeleteMessageBatchInput, ...request.Option) *sqs.DeleteMessageBatchOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.DeleteMessageBatchOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.DeleteMessageBatchInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// DeleteMessageRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) DeleteMessageRequest(_a0 *sqs.DeleteMessageInput) (*request.Request, *sqs.DeleteMessageOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.DeleteMessageInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.DeleteMessageOutput + if rf, ok := ret.Get(1).(func(*sqs.DeleteMessageInput) *sqs.DeleteMessageOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.DeleteMessageOutput) + } + } + + return r0, r1 +} + +// DeleteMessageWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) DeleteMessageWithContext(_a0 context.Context, _a1 *sqs.DeleteMessageInput, _a2 ...request.Option) (*sqs.DeleteMessageOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.DeleteMessageOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.DeleteMessageInput, ...request.Option) *sqs.DeleteMessageOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.DeleteMessageOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.DeleteMessageInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// DeleteQueue provides a mock function with given fields: _a0 +func (_m *SQSAPI) DeleteQueue(_a0 *sqs.DeleteQueueInput) (*sqs.DeleteQueueOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.DeleteQueueOutput + if rf, ok := ret.Get(0).(func(*sqs.DeleteQueueInput) *sqs.DeleteQueueOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.DeleteQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.DeleteQueueInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// DeleteQueueRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) DeleteQueueRequest(_a0 *sqs.DeleteQueueInput) (*request.Request, *sqs.DeleteQueueOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.DeleteQueueInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.DeleteQueueOutput + if rf, ok := ret.Get(1).(func(*sqs.DeleteQueueInput) *sqs.DeleteQueueOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.DeleteQueueOutput) + } + } + + return r0, r1 +} + +// DeleteQueueWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) DeleteQueueWithContext(_a0 context.Context, _a1 *sqs.DeleteQueueInput, _a2 ...request.Option) (*sqs.DeleteQueueOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.DeleteQueueOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.DeleteQueueInput, ...request.Option) *sqs.DeleteQueueOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.DeleteQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.DeleteQueueInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// GetQueueAttributes provides a mock function with given fields: _a0 +func (_m *SQSAPI) GetQueueAttributes(_a0 *sqs.GetQueueAttributesInput) (*sqs.GetQueueAttributesOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.GetQueueAttributesOutput + if rf, ok := ret.Get(0).(func(*sqs.GetQueueAttributesInput) *sqs.GetQueueAttributesOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.GetQueueAttributesOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.GetQueueAttributesInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// GetQueueAttributesRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) GetQueueAttributesRequest(_a0 *sqs.GetQueueAttributesInput) (*request.Request, *sqs.GetQueueAttributesOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.GetQueueAttributesInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.GetQueueAttributesOutput + if rf, ok := ret.Get(1).(func(*sqs.GetQueueAttributesInput) *sqs.GetQueueAttributesOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.GetQueueAttributesOutput) + } + } + + return r0, r1 +} + +// GetQueueAttributesWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) GetQueueAttributesWithContext(_a0 context.Context, _a1 *sqs.GetQueueAttributesInput, _a2 ...request.Option) (*sqs.GetQueueAttributesOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.GetQueueAttributesOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.GetQueueAttributesInput, ...request.Option) *sqs.GetQueueAttributesOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.GetQueueAttributesOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.GetQueueAttributesInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// GetQueueUrl provides a mock function with given fields: _a0 +func (_m *SQSAPI) GetQueueUrl(_a0 *sqs.GetQueueUrlInput) (*sqs.GetQueueUrlOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.GetQueueUrlOutput + if rf, ok := ret.Get(0).(func(*sqs.GetQueueUrlInput) *sqs.GetQueueUrlOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.GetQueueUrlOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.GetQueueUrlInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// GetQueueUrlRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) GetQueueUrlRequest(_a0 *sqs.GetQueueUrlInput) (*request.Request, *sqs.GetQueueUrlOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.GetQueueUrlInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.GetQueueUrlOutput + if rf, ok := ret.Get(1).(func(*sqs.GetQueueUrlInput) *sqs.GetQueueUrlOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.GetQueueUrlOutput) + } + } + + return r0, r1 +} + +// GetQueueUrlWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) GetQueueUrlWithContext(_a0 context.Context, _a1 *sqs.GetQueueUrlInput, _a2 ...request.Option) (*sqs.GetQueueUrlOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.GetQueueUrlOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.GetQueueUrlInput, ...request.Option) *sqs.GetQueueUrlOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.GetQueueUrlOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.GetQueueUrlInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ListDeadLetterSourceQueues provides a mock function with given fields: _a0 +func (_m *SQSAPI) ListDeadLetterSourceQueues(_a0 *sqs.ListDeadLetterSourceQueuesInput) (*sqs.ListDeadLetterSourceQueuesOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.ListDeadLetterSourceQueuesOutput + if rf, ok := ret.Get(0).(func(*sqs.ListDeadLetterSourceQueuesInput) *sqs.ListDeadLetterSourceQueuesOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ListDeadLetterSourceQueuesOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.ListDeadLetterSourceQueuesInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ListDeadLetterSourceQueuesPages provides a mock function with given fields: _a0, _a1 +func (_m *SQSAPI) ListDeadLetterSourceQueuesPages(_a0 *sqs.ListDeadLetterSourceQueuesInput, _a1 func(*sqs.ListDeadLetterSourceQueuesOutput, bool) bool) error { + ret := _m.Called(_a0, _a1) + + var r0 error + if rf, ok := ret.Get(0).(func(*sqs.ListDeadLetterSourceQueuesInput, func(*sqs.ListDeadLetterSourceQueuesOutput, bool) bool) error); ok { + r0 = rf(_a0, _a1) + } else { + r0 = ret.Error(0) + } + + return r0 +} + +// ListDeadLetterSourceQueuesPagesWithContext provides a mock function with given fields: _a0, _a1, _a2, _a3 +func (_m *SQSAPI) ListDeadLetterSourceQueuesPagesWithContext(_a0 context.Context, _a1 *sqs.ListDeadLetterSourceQueuesInput, _a2 func(*sqs.ListDeadLetterSourceQueuesOutput, bool) bool, _a3 ...request.Option) error { + _va := make([]interface{}, len(_a3)) + for _i := range _a3 { + _va[_i] = _a3[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1, _a2) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 error + if rf, ok := ret.Get(0).(func(context.Context, *sqs.ListDeadLetterSourceQueuesInput, func(*sqs.ListDeadLetterSourceQueuesOutput, bool) bool, ...request.Option) error); ok { + r0 = rf(_a0, _a1, _a2, _a3...) + } else { + r0 = ret.Error(0) + } + + return r0 +} + +// ListDeadLetterSourceQueuesRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) ListDeadLetterSourceQueuesRequest(_a0 *sqs.ListDeadLetterSourceQueuesInput) (*request.Request, *sqs.ListDeadLetterSourceQueuesOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.ListDeadLetterSourceQueuesInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.ListDeadLetterSourceQueuesOutput + if rf, ok := ret.Get(1).(func(*sqs.ListDeadLetterSourceQueuesInput) *sqs.ListDeadLetterSourceQueuesOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.ListDeadLetterSourceQueuesOutput) + } + } + + return r0, r1 +} + +// ListDeadLetterSourceQueuesWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) ListDeadLetterSourceQueuesWithContext(_a0 context.Context, _a1 *sqs.ListDeadLetterSourceQueuesInput, _a2 ...request.Option) (*sqs.ListDeadLetterSourceQueuesOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.ListDeadLetterSourceQueuesOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.ListDeadLetterSourceQueuesInput, ...request.Option) *sqs.ListDeadLetterSourceQueuesOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ListDeadLetterSourceQueuesOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.ListDeadLetterSourceQueuesInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ListQueueTags provides a mock function with given fields: _a0 +func (_m *SQSAPI) ListQueueTags(_a0 *sqs.ListQueueTagsInput) (*sqs.ListQueueTagsOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.ListQueueTagsOutput + if rf, ok := ret.Get(0).(func(*sqs.ListQueueTagsInput) *sqs.ListQueueTagsOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ListQueueTagsOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.ListQueueTagsInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ListQueueTagsRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) ListQueueTagsRequest(_a0 *sqs.ListQueueTagsInput) (*request.Request, *sqs.ListQueueTagsOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.ListQueueTagsInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.ListQueueTagsOutput + if rf, ok := ret.Get(1).(func(*sqs.ListQueueTagsInput) *sqs.ListQueueTagsOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.ListQueueTagsOutput) + } + } + + return r0, r1 +} + +// ListQueueTagsWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) ListQueueTagsWithContext(_a0 context.Context, _a1 *sqs.ListQueueTagsInput, _a2 ...request.Option) (*sqs.ListQueueTagsOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.ListQueueTagsOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.ListQueueTagsInput, ...request.Option) *sqs.ListQueueTagsOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ListQueueTagsOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.ListQueueTagsInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ListQueues provides a mock function with given fields: _a0 +func (_m *SQSAPI) ListQueues(_a0 *sqs.ListQueuesInput) (*sqs.ListQueuesOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.ListQueuesOutput + if rf, ok := ret.Get(0).(func(*sqs.ListQueuesInput) *sqs.ListQueuesOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ListQueuesOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.ListQueuesInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ListQueuesPages provides a mock function with given fields: _a0, _a1 +func (_m *SQSAPI) ListQueuesPages(_a0 *sqs.ListQueuesInput, _a1 func(*sqs.ListQueuesOutput, bool) bool) error { + ret := _m.Called(_a0, _a1) + + var r0 error + if rf, ok := ret.Get(0).(func(*sqs.ListQueuesInput, func(*sqs.ListQueuesOutput, bool) bool) error); ok { + r0 = rf(_a0, _a1) + } else { + r0 = ret.Error(0) + } + + return r0 +} + +// ListQueuesPagesWithContext provides a mock function with given fields: _a0, _a1, _a2, _a3 +func (_m *SQSAPI) ListQueuesPagesWithContext(_a0 context.Context, _a1 *sqs.ListQueuesInput, _a2 func(*sqs.ListQueuesOutput, bool) bool, _a3 ...request.Option) error { + _va := make([]interface{}, len(_a3)) + for _i := range _a3 { + _va[_i] = _a3[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1, _a2) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 error + if rf, ok := ret.Get(0).(func(context.Context, *sqs.ListQueuesInput, func(*sqs.ListQueuesOutput, bool) bool, ...request.Option) error); ok { + r0 = rf(_a0, _a1, _a2, _a3...) + } else { + r0 = ret.Error(0) + } + + return r0 +} + +// ListQueuesRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) ListQueuesRequest(_a0 *sqs.ListQueuesInput) (*request.Request, *sqs.ListQueuesOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.ListQueuesInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.ListQueuesOutput + if rf, ok := ret.Get(1).(func(*sqs.ListQueuesInput) *sqs.ListQueuesOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.ListQueuesOutput) + } + } + + return r0, r1 +} + +// ListQueuesWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) ListQueuesWithContext(_a0 context.Context, _a1 *sqs.ListQueuesInput, _a2 ...request.Option) (*sqs.ListQueuesOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.ListQueuesOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.ListQueuesInput, ...request.Option) *sqs.ListQueuesOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ListQueuesOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.ListQueuesInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// PurgeQueue provides a mock function with given fields: _a0 +func (_m *SQSAPI) PurgeQueue(_a0 *sqs.PurgeQueueInput) (*sqs.PurgeQueueOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.PurgeQueueOutput + if rf, ok := ret.Get(0).(func(*sqs.PurgeQueueInput) *sqs.PurgeQueueOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.PurgeQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.PurgeQueueInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// PurgeQueueRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) PurgeQueueRequest(_a0 *sqs.PurgeQueueInput) (*request.Request, *sqs.PurgeQueueOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.PurgeQueueInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.PurgeQueueOutput + if rf, ok := ret.Get(1).(func(*sqs.PurgeQueueInput) *sqs.PurgeQueueOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.PurgeQueueOutput) + } + } + + return r0, r1 +} + +// PurgeQueueWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) PurgeQueueWithContext(_a0 context.Context, _a1 *sqs.PurgeQueueInput, _a2 ...request.Option) (*sqs.PurgeQueueOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.PurgeQueueOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.PurgeQueueInput, ...request.Option) *sqs.PurgeQueueOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.PurgeQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.PurgeQueueInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ReceiveMessage provides a mock function with given fields: _a0 +func (_m *SQSAPI) ReceiveMessage(_a0 *sqs.ReceiveMessageInput) (*sqs.ReceiveMessageOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.ReceiveMessageOutput + if rf, ok := ret.Get(0).(func(*sqs.ReceiveMessageInput) *sqs.ReceiveMessageOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ReceiveMessageOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.ReceiveMessageInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// ReceiveMessageRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) ReceiveMessageRequest(_a0 *sqs.ReceiveMessageInput) (*request.Request, *sqs.ReceiveMessageOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.ReceiveMessageInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.ReceiveMessageOutput + if rf, ok := ret.Get(1).(func(*sqs.ReceiveMessageInput) *sqs.ReceiveMessageOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.ReceiveMessageOutput) + } + } + + return r0, r1 +} + +// ReceiveMessageWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) ReceiveMessageWithContext(_a0 context.Context, _a1 *sqs.ReceiveMessageInput, _a2 ...request.Option) (*sqs.ReceiveMessageOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.ReceiveMessageOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.ReceiveMessageInput, ...request.Option) *sqs.ReceiveMessageOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.ReceiveMessageOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.ReceiveMessageInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// RemovePermission provides a mock function with given fields: _a0 +func (_m *SQSAPI) RemovePermission(_a0 *sqs.RemovePermissionInput) (*sqs.RemovePermissionOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.RemovePermissionOutput + if rf, ok := ret.Get(0).(func(*sqs.RemovePermissionInput) *sqs.RemovePermissionOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.RemovePermissionOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.RemovePermissionInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// RemovePermissionRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) RemovePermissionRequest(_a0 *sqs.RemovePermissionInput) (*request.Request, *sqs.RemovePermissionOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.RemovePermissionInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.RemovePermissionOutput + if rf, ok := ret.Get(1).(func(*sqs.RemovePermissionInput) *sqs.RemovePermissionOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.RemovePermissionOutput) + } + } + + return r0, r1 +} + +// RemovePermissionWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) RemovePermissionWithContext(_a0 context.Context, _a1 *sqs.RemovePermissionInput, _a2 ...request.Option) (*sqs.RemovePermissionOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.RemovePermissionOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.RemovePermissionInput, ...request.Option) *sqs.RemovePermissionOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.RemovePermissionOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.RemovePermissionInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// SendMessage provides a mock function with given fields: _a0 +func (_m *SQSAPI) SendMessage(_a0 *sqs.SendMessageInput) (*sqs.SendMessageOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.SendMessageOutput + if rf, ok := ret.Get(0).(func(*sqs.SendMessageInput) *sqs.SendMessageOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.SendMessageOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.SendMessageInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// SendMessageBatch provides a mock function with given fields: _a0 +func (_m *SQSAPI) SendMessageBatch(_a0 *sqs.SendMessageBatchInput) (*sqs.SendMessageBatchOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.SendMessageBatchOutput + if rf, ok := ret.Get(0).(func(*sqs.SendMessageBatchInput) *sqs.SendMessageBatchOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.SendMessageBatchOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.SendMessageBatchInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// SendMessageBatchRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) SendMessageBatchRequest(_a0 *sqs.SendMessageBatchInput) (*request.Request, *sqs.SendMessageBatchOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.SendMessageBatchInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.SendMessageBatchOutput + if rf, ok := ret.Get(1).(func(*sqs.SendMessageBatchInput) *sqs.SendMessageBatchOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.SendMessageBatchOutput) + } + } + + return r0, r1 +} + +// SendMessageBatchWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) SendMessageBatchWithContext(_a0 context.Context, _a1 *sqs.SendMessageBatchInput, _a2 ...request.Option) (*sqs.SendMessageBatchOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.SendMessageBatchOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.SendMessageBatchInput, ...request.Option) *sqs.SendMessageBatchOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.SendMessageBatchOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.SendMessageBatchInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// SendMessageRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) SendMessageRequest(_a0 *sqs.SendMessageInput) (*request.Request, *sqs.SendMessageOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.SendMessageInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.SendMessageOutput + if rf, ok := ret.Get(1).(func(*sqs.SendMessageInput) *sqs.SendMessageOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.SendMessageOutput) + } + } + + return r0, r1 +} + +// SendMessageWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) SendMessageWithContext(_a0 context.Context, _a1 *sqs.SendMessageInput, _a2 ...request.Option) (*sqs.SendMessageOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.SendMessageOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.SendMessageInput, ...request.Option) *sqs.SendMessageOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.SendMessageOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.SendMessageInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// SetQueueAttributes provides a mock function with given fields: _a0 +func (_m *SQSAPI) SetQueueAttributes(_a0 *sqs.SetQueueAttributesInput) (*sqs.SetQueueAttributesOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.SetQueueAttributesOutput + if rf, ok := ret.Get(0).(func(*sqs.SetQueueAttributesInput) *sqs.SetQueueAttributesOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.SetQueueAttributesOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.SetQueueAttributesInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// SetQueueAttributesRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) SetQueueAttributesRequest(_a0 *sqs.SetQueueAttributesInput) (*request.Request, *sqs.SetQueueAttributesOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.SetQueueAttributesInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.SetQueueAttributesOutput + if rf, ok := ret.Get(1).(func(*sqs.SetQueueAttributesInput) *sqs.SetQueueAttributesOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.SetQueueAttributesOutput) + } + } + + return r0, r1 +} + +// SetQueueAttributesWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) SetQueueAttributesWithContext(_a0 context.Context, _a1 *sqs.SetQueueAttributesInput, _a2 ...request.Option) (*sqs.SetQueueAttributesOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.SetQueueAttributesOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.SetQueueAttributesInput, ...request.Option) *sqs.SetQueueAttributesOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.SetQueueAttributesOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.SetQueueAttributesInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// TagQueue provides a mock function with given fields: _a0 +func (_m *SQSAPI) TagQueue(_a0 *sqs.TagQueueInput) (*sqs.TagQueueOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.TagQueueOutput + if rf, ok := ret.Get(0).(func(*sqs.TagQueueInput) *sqs.TagQueueOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.TagQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.TagQueueInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// TagQueueRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) TagQueueRequest(_a0 *sqs.TagQueueInput) (*request.Request, *sqs.TagQueueOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.TagQueueInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.TagQueueOutput + if rf, ok := ret.Get(1).(func(*sqs.TagQueueInput) *sqs.TagQueueOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.TagQueueOutput) + } + } + + return r0, r1 +} + +// TagQueueWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) TagQueueWithContext(_a0 context.Context, _a1 *sqs.TagQueueInput, _a2 ...request.Option) (*sqs.TagQueueOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.TagQueueOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.TagQueueInput, ...request.Option) *sqs.TagQueueOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.TagQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.TagQueueInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// UntagQueue provides a mock function with given fields: _a0 +func (_m *SQSAPI) UntagQueue(_a0 *sqs.UntagQueueInput) (*sqs.UntagQueueOutput, error) { + ret := _m.Called(_a0) + + var r0 *sqs.UntagQueueOutput + if rf, ok := ret.Get(0).(func(*sqs.UntagQueueInput) *sqs.UntagQueueOutput); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.UntagQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(*sqs.UntagQueueInput) error); ok { + r1 = rf(_a0) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// UntagQueueRequest provides a mock function with given fields: _a0 +func (_m *SQSAPI) UntagQueueRequest(_a0 *sqs.UntagQueueInput) (*request.Request, *sqs.UntagQueueOutput) { + ret := _m.Called(_a0) + + var r0 *request.Request + if rf, ok := ret.Get(0).(func(*sqs.UntagQueueInput) *request.Request); ok { + r0 = rf(_a0) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*request.Request) + } + } + + var r1 *sqs.UntagQueueOutput + if rf, ok := ret.Get(1).(func(*sqs.UntagQueueInput) *sqs.UntagQueueOutput); ok { + r1 = rf(_a0) + } else { + if ret.Get(1) != nil { + r1 = ret.Get(1).(*sqs.UntagQueueOutput) + } + } + + return r0, r1 +} + +// UntagQueueWithContext provides a mock function with given fields: _a0, _a1, _a2 +func (_m *SQSAPI) UntagQueueWithContext(_a0 context.Context, _a1 *sqs.UntagQueueInput, _a2 ...request.Option) (*sqs.UntagQueueOutput, error) { + _va := make([]interface{}, len(_a2)) + for _i := range _a2 { + _va[_i] = _a2[_i] + } + var _ca []interface{} + _ca = append(_ca, _a0, _a1) + _ca = append(_ca, _va...) + ret := _m.Called(_ca...) + + var r0 *sqs.UntagQueueOutput + if rf, ok := ret.Get(0).(func(context.Context, *sqs.UntagQueueInput, ...request.Option) *sqs.UntagQueueOutput); ok { + r0 = rf(_a0, _a1, _a2...) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*sqs.UntagQueueOutput) + } + } + + var r1 error + if rf, ok := ret.Get(1).(func(context.Context, *sqs.UntagQueueInput, ...request.Option) error); ok { + r1 = rf(_a0, _a1, _a2...) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +type mockConstructorTestingTNewSQSAPI interface { + mock.TestingT + Cleanup(func()) +} + +// NewSQSAPI creates a new instance of SQSAPI. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. +func NewSQSAPI(t mockConstructorTestingTNewSQSAPI) *SQSAPI { + mock := &SQSAPI{} + mock.Mock.Test(t) + + t.Cleanup(func() { mock.AssertExpectations(t) }) + + return mock +} diff --git a/idk/ingest.go b/idk/ingest.go index 85dc85793..38b68d5af 100644 --- a/idk/ingest.go +++ b/idk/ingest.go @@ -191,6 +191,7 @@ func init() { func (m *Main) PilosaClient() *pilosaclient.Client { return m.client } func (m *Main) Log() logger.Logger { return m.log } +func (m *Main) SetLog(log logger.Logger) { m.log = log } func NewMain() *Main { fmt.Fprintf(os.Stderr, "Molecula Consumer %s, build time %s\n", Version, BuildTime) diff --git a/idk/kinesis/cmd.go b/idk/kinesis/cmd.go index 2e070675a..bb9b9bded 100644 --- a/idk/kinesis/cmd.go +++ b/idk/kinesis/cmd.go @@ -3,10 +3,15 @@ package kinesis import ( "time" - "github.com/molecula/featurebase/v3/idk" "github.com/pkg/errors" + + "github.com/molecula/featurebase/v3/idk" ) +// Main is the holder of all configurations for a Kinesis stream consumer. +// +// Along with the additional configuration fields, kinesis.Main also gains all +// fields and methods from idk.Main via composition. type Main struct { idk.Main `flag:"!embed"` Timeout time.Duration `help:"Time to wait for more records from Kinesis before flushing a batch. 0 to disable."` @@ -16,8 +21,28 @@ type Main struct { StreamName string `help:"Name of AWS Kinesis stream to consume records from."` OffsetsPath string `help:"Path where the offsets file will be written. May be a path on the local filesystem, or an S3 URI."` AWSProfile string `help:"Name of AWS profile to use. Alternatively, use environment variable AWS_PROFILE."` + ErrorQueueName string `help:"SQS queue name to send error and panic/runtime errors to."` } +// NewMain returns a new instance of a Kinesis stream consumer configuration object. +// +// It specifies a callback NewSource that can be invoked to create a kinesis.Source object. +// This callback implicitly initializes an AWS session and uses that session to initialize +// clients to the following AWS resources: S3, Kinesis, and SQS. Client creation happens regardless +// of configuration. (ex: OffsetsPath and Header are local paths -> S3 client is created.) +// +// The default BatchSize is 20000 and Concurrency is 1. Any Concurrency value > 1 is NOT supported. +// These values are set on the returned kinesis.Main instance. +// +// The Logger instance on the kinesis.Source is always decorated when NewSource is invoked. +// Assuming no errors occur during AWS client initialization, the decorated Logger instance is +// propagated back to the kinesis.Main so that callers that configured it can also emit errors and +// panics to the SQS queue specified by ErrorQueueName. The behavior of the wrapped Logger depends on +// a non-empty ErrorQueueName, the existence of an SQS queue instance in AWSRegion with that name +// the StreamName field being of a particular format 'PREFIX'-VALID_UUID, and if a valid SQS queue URL +// can be resolved at the time of Logger initialization. If any of these are false, the error +// emission to an SQS queue functionality is not activated and the Logger instance behaves identically +// to its wrapped Logger and emits a warning to the caller that errors are not propagated to SQS. func NewMain() *Main { m := &Main{ Main: *idk.NewMain(), @@ -28,9 +53,9 @@ func NewMain() *Main { m.OffsetMode = true m.Main.Namespace = "ingester_kinesis" m.Main.Pprof = "" // don't initialize pprof until we actually use it in tests + m.NewSource = func() (idk.Source, error) { source := NewSource() - source.Log = m.Main.Log() source.Timeout = m.Timeout source.Header = m.Header source.AWSRegion = m.AWSRegion @@ -39,10 +64,21 @@ func NewMain() *Main { source.OffsetsPath = m.OffsetsPath source.AWSProfile = m.AWSProfile + // This Logger instance is wrapped in `Open` -> `initAWS`. + source.Log = m.Main.Log() + source.ErrorQueueName = m.ErrorQueueName + err := source.Open() if err != nil { return nil, errors.Wrap(err, "opening source") } + + // `Open` succeeded -> AWS resources successfully initialized -> + // Logger instance was successfully wrapped. Now assign the wrapped + // Logger instance back to main so executables invoking this + // (ex: `molecula-consumer-kinesis`) will propagate errors and + // panics on failure using the wrapped Logger instance. + m.Main.SetLog(source.Log) return source, nil } return m diff --git a/idk/kinesis/logger.go b/idk/kinesis/logger.go new file mode 100644 index 000000000..88f9cc9aa --- /dev/null +++ b/idk/kinesis/logger.go @@ -0,0 +1,234 @@ +package kinesis + +import ( + "encoding/base64" + "encoding/json" + "fmt" + "strings" + "time" + + "github.com/aws/aws-sdk-go/aws" + "github.com/aws/aws-sdk-go/service/sqs" + "github.com/aws/aws-sdk-go/service/sqs/sqsiface" + "github.com/google/uuid" + "github.com/pkg/errors" + + "github.com/molecula/featurebase/v3/logger" +) + +// ErrorType signifies the type of error encountered. +type ErrorType string + +// These are the currently supported ErrorType values that can be emitted +// to an ErrorStore. +const ( + RecoverableErrorType = ErrorType("Error") + PanicErrorType = ErrorType("Panic") // Runtime errors. +) + +// ErrorStore is an abstraction over a resource like external storage, a database, +// a queue, etc. that can receive and store error messages. +type ErrorStore interface { + // Available checks if the backing resource can receive error + // messages via Push. + Available() bool + + // Push emits an error message of type ErrorType to the backing resource. + // + // The caller should NOT assume that Available is called implicitly + // to check for ErrorStore availability. + Push(ErrorType, string, logger.Logger) error +} + +// SinkErrorPayload contains all data about a single IDK error message that will be +// emitted to an ErrorStore; includes a timestamp and a valid sink ID. +// +// This payload is not meant to be used outside the context of error propogation to +// an ErrorStore. +type SinkErrorPayload struct { + SinkId string `json:"sink_id"` + ErrorType ErrorType `json:"error_type"` + ErrorMessage string `json:"error_msg"` + Timestamp string `json:"time"` +} + +// SinkErrorQueue is an ErrorStore implementation that uses an SQS queue as its backing +// resource to emit error and panic messages to. +// +// It also maps 1-to-1 to a Kinesis stream via a unique sink ID. +type SinkErrorQueue struct { + sinkId string + name string + url string + queue sqsiface.SQSAPI +} + +// NewSinkErrorQueue attempts to construct a SinkErrorQueue instance from an AWS SQS client, +// a queue name, and a sink ID. +// +// On success, callers can assume that a backing SQS queue resource exists and is fully initialized. +// +// Returns nil and an SQS error if an SQS queue URL cannot be resolved from the queue name and/or +// the AWS SQS client. +// +// This method assumes the sink ID argument is valid. +func NewSinkErrorQueue(queue sqsiface.SQSAPI, queueName, sinkId string) (*SinkErrorQueue, error) { + input := &sqs.GetQueueUrlInput{QueueName: &queueName} + output, err := queue.GetQueueUrl(input) + if err != nil { + return nil, err + } + + return &SinkErrorQueue{sinkId, queueName, *output.QueueUrl, queue}, nil +} + +// SinkErrorQueueFrom always constructs a SinkErrorQueue instance from an AWS SQS client and a +// kinesis.Source. +// +// Unlike NewSinkErrorQueue, this does NOT return an error if a queue URL cannot be resolved +// from the queue name and/or the AWS SQS client. Instead it will collapse to a SinkErrorQueue +// instance with a backing SQS resource that is ALWAYS unavailable. Attempting to invoke Push +// on this instance will not result in an error; instead it will just emit a warning that no +// backing SQS resource could be written to. +// +// This does check if the sink ID has a valid form: 'PREFIX'-VALID_UUID. If not, this collapses +// to a SinkErrorQueue instance that is ALWAYS unavailable. +func SinkErrorQueueFrom(queue sqsiface.SQSAPI, source *Source) *SinkErrorQueue { + // The below failure conditions are handled by collapsing to a no-op ErrorStore.Push implementation. + // - A missing SQS queue name. + // - An invalid sink ID (i.e. not a valid UUID); valid sink ID: "PREFIX"-UUID. + // - Unable to resolve queue URL from queue name. + // + // Means downstream ErrorStreamLogger behaves identical to its embedded Logger. + if source.ErrorQueueName == "" { + return &SinkErrorQueue{} + } + + sinkId := strings.Join(strings.Split(source.StreamName, "-")[1:], "-") + _, err := uuid.Parse(sinkId) + if err != nil { + // Keep invalid sink ID around in case something downstream wants to log. + return &SinkErrorQueue{sinkId, source.ErrorQueueName, "", nil} + } + + sinkErrorQueue, err := NewSinkErrorQueue(queue, source.ErrorQueueName, sinkId) + if err != nil { + return &SinkErrorQueue{sinkId, source.ErrorQueueName, "", nil} + } + + return sinkErrorQueue +} + +// Available checks that a valid SQS queue resource exists. +func (seq *SinkErrorQueue) Available() bool { + return seq.url != "" && seq.queue != nil +} + +// Push attempts to emit a single error message of type ErrorType to a SQS queue resource. +// +// If the backing SQS queue resource is not available, this function is a no-op and does NOT +// return an error. Instead it emits a warning to the logger.Logger instance specified by the log +// argument. +// +// The warning can be ignored entirely by passing a nil log argument. +func (seq *SinkErrorQueue) Push(errorType ErrorType, message string, log logger.Logger) error { + if !seq.Available() { + if log != nil { + log.Warnf("Not pushing errors to an SQS queue='%+v' due to unavailability.", seq) + } + return nil + } + + payload := SinkErrorPayload{ + SinkId: seq.sinkId, + ErrorType: errorType, + ErrorMessage: message, + Timestamp: time.Now().Format(time.RFC3339), + } + + payloadBytes, err := json.Marshal(&payload) + if err != nil { + msgTemplate := "Unable to marshal payload='%+v' to send message='%s' to queue='%+v'." + return errors.Wrap(err, fmt.Sprintf(msgTemplate, payload, message, seq.queue)) + } + + encodedPayload := base64.URLEncoding.EncodeToString(payloadBytes) + input := &sqs.SendMessageInput{ + DelaySeconds: aws.Int64(10), + MessageBody: aws.String(encodedPayload), + QueueUrl: &seq.url, + } + + _, err = seq.queue.SendMessage(input) + if err != nil { + msgTemplate := "Unable to send message='%s' to queue='%+v' with input='%+v'." + return errors.Wrap(err, fmt.Sprintf(msgTemplate, message, seq.queue, input)) + } + return nil +} + +// ErrorStreamLogger is a logger.Logger implementation that decorates a base logger.Logger instance +// and emits error and panic messages to an ErrorStore. +// +// All other log levels delegate to the base logger.Logger implementation. +type ErrorStreamLogger struct { + base logger.Logger + store ErrorStore +} + +// NewErrorStreamLogger constructs an ErrorStreamLogger from a logger.Logger and an ErrorStore. +func NewErrorStreamLogger(base logger.Logger, store ErrorStore) *ErrorStreamLogger { + return &ErrorStreamLogger{base, store} +} + +// Printf just delegates to the wrapped/base logger.Logger's Printf implementation. +func (esl *ErrorStreamLogger) Printf(format string, v ...interface{}) { + esl.base.Printf(format, v...) +} + +// Debugf just delegates to the wrapped/base logger.Logger's Debugf implementation. +func (esl *ErrorStreamLogger) Debugf(format string, v ...interface{}) { + esl.base.Debugf(format, v...) +} + +// Infof just delegates to the wrapped/base logger.Logger's Infof implementation. +func (esl *ErrorStreamLogger) Infof(format string, v ...interface{}) { + esl.base.Infof(format, v...) +} + +// Warnf just delegates to the wrapped/base logger.Logger's Warnf implementation. +func (esl *ErrorStreamLogger) Warnf(format string, v ...interface{}) { + esl.base.Warnf(format, v...) +} + +// Errorf delegates to the wrapped/base logger.Logger's Errorf implementation and additionally/ +// pushes an error message with RecoverableErrorType ErrorType to the ErrorStore. +// +// If an error occurs during a Push to the ErrorStore, the error is logged to the wrapped/base +// logger.Logger using Errorf. +func (esl *ErrorStreamLogger) Errorf(format string, v ...interface{}) { + esl.base.Errorf(format, v...) + + err := esl.store.Push(RecoverableErrorType, fmt.Sprintf(format, v...), esl) + if err != nil { + errMsg := fmt.Sprintf("Failed during push to store='%+v'", esl.store) + esl.base.Errorf(errors.Wrap(err, errMsg).Error()) + } +} + +// Panicf delegates to the wrapped/base logger.Logger's Panicf implementation and additionally +// pushes an error message with PanicErrorType ErrorType to the ErrorStore. +// +// If an error occurs during a Push to the ErrorStore, the error is logged to the wrapped/base +// logger.Logger using Errorf NOT Panicf. +func (esl *ErrorStreamLogger) Panicf(format string, v ...interface{}) { + esl.base.Panicf(format, v...) + + err := esl.store.Push(PanicErrorType, fmt.Sprintf(format, v...), esl) + if err != nil { + // A (recoverable) error occurred during a push to the store, so + // log that error using `Errorf` not `Panicf`. + errMsg := fmt.Sprintf("Failed during push to store='%+v'", esl.store) + esl.base.Errorf(errors.Wrap(err, errMsg).Error()) + } +} diff --git a/idk/kinesis/logger_test.go b/idk/kinesis/logger_test.go new file mode 100644 index 000000000..074621c05 --- /dev/null +++ b/idk/kinesis/logger_test.go @@ -0,0 +1,540 @@ +package kinesis + +import ( + "encoding/base64" + "encoding/json" + "fmt" + "strings" + "testing" + "time" + + "github.com/aws/aws-sdk-go/aws" + "github.com/aws/aws-sdk-go/service/sqs" + "github.com/pkg/errors" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" + + "github.com/molecula/featurebase/v3/idk/idktest/mocks" + "github.com/molecula/featurebase/v3/logger" +) + +// In-memory featurebase/logger.Logger implementation that stores by level +// {"PRINTF", "DEBUG", "INFO", "WARN", "ERROR", "PANIC"} useful for unit tests. +type MapStashLogger struct { + messages map[string][]string +} + +func NewMapStashLogger() *MapStashLogger { + messages := make(map[string][]string) + messages["PRINTF"] = []string{} + messages["DEBUG"] = []string{} + messages["INFO"] = []string{} + messages["WARN"] = []string{} + messages["ERROR"] = []string{} + messages["PANIC"] = []string{} + return &MapStashLogger{messages} +} + +func (msl *MapStashLogger) Printf(format string, v ...interface{}) { + msl.messages["PRINTF"] = append(msl.messages["PRINTF"], fmt.Sprintf(format, v...)) +} + +func (msl *MapStashLogger) Debugf(format string, v ...interface{}) { + msl.messages["DEBUG"] = append(msl.messages["DEBUG"], fmt.Sprintf(format, v...)) +} + +func (msl *MapStashLogger) Infof(format string, v ...interface{}) { + msl.messages["INFO"] = append(msl.messages["INFO"], fmt.Sprintf(format, v...)) +} + +func (msl *MapStashLogger) Warnf(format string, v ...interface{}) { + msl.messages["WARN"] = append(msl.messages["WARN"], fmt.Sprintf(format, v...)) +} + +func (msl *MapStashLogger) Errorf(format string, v ...interface{}) { + msl.messages["ERROR"] = append(msl.messages["ERROR"], fmt.Sprintf(format, v...)) +} + +func (msl *MapStashLogger) Panicf(format string, v ...interface{}) { + msl.messages["PANIC"] = append(msl.messages["PANIC"], fmt.Sprintf(format, v...)) +} + +// In-memory error store implementation to use in unit tests. +type MapStashErrorStore struct { + storage map[ErrorType][]string + available bool + simulatedError error +} + +func NewMapStashErrorStore(available bool, err error) *MapStashErrorStore { + return &MapStashErrorStore{make(map[ErrorType][]string), available, err} +} + +func (mses *MapStashErrorStore) Available() bool { + return mses.available +} + +func (mses *MapStashErrorStore) Push(errorType ErrorType, message string, log logger.Logger) error { + if !mses.Available() { + return nil + } + + if mses.simulatedError != nil { + return mses.simulatedError + } + + mses.storage[errorType] = append(mses.storage[errorType], message) + return nil +} + +func TestNewSinkErrorQueueSuccess(t *testing.T) { + mockSQS := &mocks.SQSAPI{} + + mockGetQueueUrlFn := func(input *sqs.GetQueueUrlInput) *sqs.GetQueueUrlOutput { + name := *input.QueueName + url := fmt.Sprintf("https://unit-test.queue.%s.url", name) + return &sqs.GetQueueUrlOutput{QueueUrl: aws.String(url)} + } + + mockSQS.On("GetQueueUrl", mock.MatchedBy(func(input *sqs.GetQueueUrlInput) bool { + return input.QueueName != nil && *input.QueueName != "" + })).Return(mockGetQueueUrlFn, nil) + + queue, err := NewSinkErrorQueue(mockSQS, "dummy-123", "a-b-c") + + assert.Nil(t, err) + assert.Equal(t, "a-b-c", queue.sinkId) + assert.Equal(t, "dummy-123", queue.name) + assert.Equal(t, "https://unit-test.queue.dummy-123.url", queue.url) + assert.Same(t, mockSQS, queue.queue) +} + +func TestNewSinkErrorQueueFailSQSGetQueueUrlErrored(t *testing.T) { + mockSQS := &mocks.SQSAPI{} + + mockGetQueueUrlFn := func(input *sqs.GetQueueUrlInput) *sqs.GetQueueUrlOutput { + // Value ignored by the caller on failure. + return &sqs.GetQueueUrlOutput{} + } + + errMsg := "Could not retrieve queue URL." + mockSQS.On("GetQueueUrl", mock.MatchedBy(func(input *sqs.GetQueueUrlInput) bool { + return true + })).Return(mockGetQueueUrlFn, errors.New(errMsg)) + + queue, err := NewSinkErrorQueue(mockSQS, "dummy-123", "a-b-c") + + assert.Nil(t, queue) + assert.NotNil(t, err) + assert.Equal(t, errMsg, err.Error()) +} + +func TestSinkErrorQueueFromSuccess(t *testing.T) { + mockSQS := &mocks.SQSAPI{} + + mockGetQueueUrlFn := func(input *sqs.GetQueueUrlInput) *sqs.GetQueueUrlOutput { + name := *input.QueueName + url := fmt.Sprintf("https://unit-test.queue.%s.url", name) + return &sqs.GetQueueUrlOutput{QueueUrl: aws.String(url)} + } + + mockSQS.On("GetQueueUrl", mock.MatchedBy(func(input *sqs.GetQueueUrlInput) bool { + return input.QueueName != nil && *input.QueueName != "" + })).Return(mockGetQueueUrlFn, nil) + + validUuid := "0f2af1d9-52db-4a1c-bc75-554d84b01850" + streamName := fmt.Sprintf("sink-%s", validUuid) // Structure mimics how cloud ECS instance names Kinesis stream. + source := &Source{ + ErrorQueueName: "my-queue-01234", + StreamName: streamName, + } + + var queue *SinkErrorQueue + queue = SinkErrorQueueFrom(mockSQS, source) + + assert.Equal(t, validUuid, queue.sinkId) + assert.Equal(t, "my-queue-01234", queue.name) + assert.Equal(t, "https://unit-test.queue.my-queue-01234.url", queue.url) + assert.Same(t, mockSQS, queue.queue) +} + +func TestSinkErrorQueueFromFailMissingQueueName(t *testing.T) { + // Missing queue name on `Source` object creates a failure condition before SQS instance is even inspected. + mockSQS := &mocks.SQSAPI{} + + validUuid := "0f2af1d9-52db-4a1c-bc75-554d84b01850" + streamName := fmt.Sprintf("sink-%s", validUuid) // Structure mimics how cloud ECS instance names its Kinesis stream. + source := &Source{StreamName: streamName} + + var queue *SinkErrorQueue + queue = SinkErrorQueueFrom(mockSQS, source) + + assert.Equal(t, "", queue.sinkId) + assert.Equal(t, "", queue.name) + assert.Equal(t, "", queue.url) + assert.Nil(t, queue.queue) +} + +func TestSinkErrorQueueFromFailMalformedSinkId(t *testing.T) { + mockSQS := &mocks.SQSAPI{} + + mockGetQueueUrlFn := func(input *sqs.GetQueueUrlInput) *sqs.GetQueueUrlOutput { + // Value ignored by the caller on failure. + return &sqs.GetQueueUrlOutput{} + } + + mockSQS.On("GetQueueUrl", mock.MatchedBy(func(input *sqs.GetQueueUrlInput) bool { + return true + })).Return(mockGetQueueUrlFn, nil) + + malformedUuid := "28hsdfas636bd" + streamName := fmt.Sprintf("sink-%s", malformedUuid) + source := &Source{ + ErrorQueueName: "valid-queue-90123", + StreamName: streamName, + } + + var queue *SinkErrorQueue + queue = SinkErrorQueueFrom(mockSQS, source) + + assert.Equal(t, malformedUuid, queue.sinkId) // Keeps invalid sink ID around for downstream logging purposes. + assert.Equal(t, "valid-queue-90123", queue.name) + assert.Equal(t, "", queue.url) // Did not reach a point where `sqs.GetQueueUrl` is even called. + assert.Nil(t, queue.queue) +} + +func TestSinkErrorQueueFromFailEmptyStreamNameDoesNotCausePanic(t *testing.T) { + mockSQS := &mocks.SQSAPI{} + + mockGetQueueUrlFn := func(input *sqs.GetQueueUrlInput) *sqs.GetQueueUrlOutput { + // Value ignored by the caller on failure. + return &sqs.GetQueueUrlOutput{} + } + + mockSQS.On("GetQueueUrl", mock.MatchedBy(func(input *sqs.GetQueueUrlInput) bool { + return true + })).Return(mockGetQueueUrlFn, nil) + + source := &Source{ + ErrorQueueName: "valid-queue-88888888888", + StreamName: "", + } + + var queue *SinkErrorQueue + queue = SinkErrorQueueFrom(mockSQS, source) + + assert.Equal(t, "", queue.sinkId) // Keeps invalid sink ID around for downstream logging purposes. + assert.Equal(t, "valid-queue-88888888888", queue.name) + assert.Equal(t, "", queue.url) // Did not reach a point where `sqs.GetQueueUrl` is even called. + assert.Nil(t, queue.queue) +} + +func TestSinkErrorQueueFromFailSQSGetQueueUrlErrored(t *testing.T) { + mockSQS := &mocks.SQSAPI{} + + mockGetQueueUrlFn := func(input *sqs.GetQueueUrlInput) *sqs.GetQueueUrlOutput { + // Value ignored by the caller on failure. + return &sqs.GetQueueUrlOutput{} + } + + errMsg := "This error message isn't propagated but is handled within the method." + mockSQS.On("GetQueueUrl", mock.MatchedBy(func(input *sqs.GetQueueUrlInput) bool { + return true + })).Return(mockGetQueueUrlFn, errors.New(errMsg)) + + validUuid := "fa6f1631-7c29-4023-95ac-48a2f59fae9b" + streamName := fmt.Sprintf("sink-%s", validUuid) // Structure mimics how cloud ECS instance names Kinesis stream. + source := &Source{ + ErrorQueueName: "queue-that-we-stole-962463", + StreamName: streamName, + } + + var queue *SinkErrorQueue + queue = SinkErrorQueueFrom(mockSQS, source) + + assert.Equal(t, validUuid, queue.sinkId) // Keeps valid sink ID around for downstream logging purposes. + assert.Equal(t, "queue-that-we-stole-962463", queue.name) + assert.Equal(t, "", queue.url) // Could not resolve queue URL, so this remains empty. + assert.Nil(t, queue.queue) +} + +func TestSinkErrorQueueAvailable(t *testing.T) { + emptySeq := &SinkErrorQueue{} + assert.False(t, emptySeq.Available()) + + // Only depends on URL being non-empty and a non-nil reference to an SQS queue. + mockSQS := &mocks.SQSAPI{} + seq := &SinkErrorQueue{ + url: "asdfadfasdf2332798", + queue: mockSQS, + } + assert.True(t, seq.Available()) +} + +func TestSinkErrorQueuePushWhenQueueNotAvailableDoesNotThrowErrorAndNoOps(t *testing.T) { + emptySeq := &SinkErrorQueue{} + + // Logger is nil. + err := emptySeq.Push(RecoverableErrorType, "werqw zxcw7228323974", nil) + assert.Nil(t, err) + + // Logger is not nil; check that a warning is issued using the Logger instance. + logger := NewMapStashLogger() + err2 := emptySeq.Push(RecoverableErrorType, "werqw zxcw7228323974", logger) + assert.Nil(t, err2) + assert.Equal(t, 1, len(logger.messages["WARN"])) + + // Check warning is emitted to the logger. + assert.True(t, strings.HasPrefix(logger.messages["WARN"][0], "Not pushing errors to an SQS queue=")) +} + +func TestSinkErrorQueuePushSuccess(t *testing.T) { + mockSQS := &mocks.SQSAPI{} + + // Use variables defined outside the scope of a closure to save the message body and queueURL. + var actualBody string + var actualQueueUrl string + mockSendMessageFn := func(input *sqs.SendMessageInput) *sqs.SendMessageOutput { + actualBody = *input.MessageBody + actualQueueUrl = *input.QueueUrl + return &sqs.SendMessageOutput{MessageId: aws.String(actualBody[:5])} + } + + mockSQS.On("SendMessage", mock.MatchedBy(func(input *sqs.SendMessageInput) bool { + return input.QueueUrl != nil && *input.QueueUrl != "" + })).Return(mockSendMessageFn, nil) + + seq := &SinkErrorQueue{ + sinkId: "xxx-000", + name: "pushover", + url: "https://unit-test.queue.push.url", + queue: mockSQS, + } + + errMsg := "079 cxmn, 198sfakjl" + err := seq.Push(PanicErrorType, errMsg, nil) + assert.Nil(t, err) + + // Verify queue URL. + assert.Equal(t, seq.url, actualQueueUrl) + + // Now check message body. + decoded, err := base64.URLEncoding.DecodeString(actualBody) + assert.Nil(t, err, "Error while base64 URL decoding.") + + payload := SinkErrorPayload{} + err = json.Unmarshal(decoded, &payload) + assert.Nil(t, err, "Error while de-serializing JSON into SinkErrorPayload object.") // An error de-serializing JSON fails the unit test. + + assert.Equal(t, "xxx-000", payload.SinkId) + assert.Equal(t, "Panic", string(payload.ErrorType)) + assert.Equal(t, errMsg, payload.ErrorMessage) + assert.GreaterOrEqual(t, time.Now().Format(time.RFC3339), payload.Timestamp) // Now's timestamp should be later than timestamp of sent message. +} + +func TestNewErrorStreamLogger(t *testing.T) { + baseLogger := NewMapStashLogger() + store := NewMapStashErrorStore(true, nil) + + logger := NewErrorStreamLogger(baseLogger, store) + assert.Same(t, baseLogger, logger.base) + assert.Same(t, store, logger.store) +} + +func TestErrorStreamLoggerPrintf(t *testing.T) { + baseLogger := NewMapStashLogger() + store := NewMapStashErrorStore(true, nil) + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Printf("printf:%d", 3) + + assert.Equal(t, 1, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 0, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 0, len(baseLogger.messages["INFO"])) + assert.Equal(t, 0, len(baseLogger.messages["WARN"])) + assert.Equal(t, 0, len(baseLogger.messages["ERROR"])) + assert.Equal(t, 0, len(baseLogger.messages["PANIC"])) + assert.Equal(t, "printf:3", baseLogger.messages["PRINTF"][0]) + + assert.Equal(t, 0, len(store.storage[RecoverableErrorType])) + assert.Equal(t, 0, len(store.storage[PanicErrorType])) +} + +func TestErrorStreamLoggerDebugf(t *testing.T) { + baseLogger := NewMapStashLogger() + store := NewMapStashErrorStore(true, nil) + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Debugf("debug::%d", 4) + + assert.Equal(t, 0, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 1, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 0, len(baseLogger.messages["INFO"])) + assert.Equal(t, 0, len(baseLogger.messages["WARN"])) + assert.Equal(t, 0, len(baseLogger.messages["ERROR"])) + assert.Equal(t, 0, len(baseLogger.messages["PANIC"])) + assert.Equal(t, "debug::4", baseLogger.messages["DEBUG"][0]) + + assert.Equal(t, 0, len(store.storage[RecoverableErrorType])) + assert.Equal(t, 0, len(store.storage[PanicErrorType])) +} + +func TestErrorStreamLoggerInfof(t *testing.T) { + baseLogger := NewMapStashLogger() + store := NewMapStashErrorStore(true, nil) + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Infof("info:::%d", 5) + + assert.Equal(t, 0, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 0, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 1, len(baseLogger.messages["INFO"])) + assert.Equal(t, 0, len(baseLogger.messages["WARN"])) + assert.Equal(t, 0, len(baseLogger.messages["ERROR"])) + assert.Equal(t, 0, len(baseLogger.messages["PANIC"])) + assert.Equal(t, "info:::5", baseLogger.messages["INFO"][0]) + + assert.Equal(t, 0, len(store.storage[RecoverableErrorType])) + assert.Equal(t, 0, len(store.storage[PanicErrorType])) +} + +func TestErrorStreamLoggerWarnf(t *testing.T) { + baseLogger := NewMapStashLogger() + store := NewMapStashErrorStore(true, nil) + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Warnf("warn::::%d", 6) + + assert.Equal(t, 0, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 0, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 0, len(baseLogger.messages["INFO"])) + assert.Equal(t, 1, len(baseLogger.messages["WARN"])) + assert.Equal(t, 0, len(baseLogger.messages["ERROR"])) + assert.Equal(t, 0, len(baseLogger.messages["PANIC"])) + assert.Equal(t, "warn::::6", baseLogger.messages["WARN"][0]) + + assert.Equal(t, 0, len(store.storage[RecoverableErrorType])) + assert.Equal(t, 0, len(store.storage[PanicErrorType])) +} + +func TestErrorStreamLoggerErrorfWithAvailableQueue(t *testing.T) { + baseLogger := NewMapStashLogger() + store := NewMapStashErrorStore(true, nil) // Queue is available so a push actually occurs. + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Errorf("error:::::%d", 7) + + assert.Equal(t, 0, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 0, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 0, len(baseLogger.messages["INFO"])) + assert.Equal(t, 0, len(baseLogger.messages["WARN"])) + assert.Equal(t, 1, len(baseLogger.messages["ERROR"])) + assert.Equal(t, 0, len(baseLogger.messages["PANIC"])) + assert.Equal(t, "error:::::7", baseLogger.messages["ERROR"][0]) + + assert.Equal(t, 1, len(store.storage[RecoverableErrorType])) // Verifies a push occurred. + assert.Equal(t, 0, len(store.storage[PanicErrorType])) + assert.Equal(t, "error:::::7", store.storage[RecoverableErrorType][0]) +} + +func TestErrorStreamLoggerErrorfNoAvailableQueue(t *testing.T) { + baseLogger := NewMapStashLogger() + store := NewMapStashErrorStore(false, nil) // Queue is unavailable. + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Errorf("error:::::%d%s", 7, "a") + + assert.Equal(t, 0, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 0, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 0, len(baseLogger.messages["INFO"])) + assert.Equal(t, 0, len(baseLogger.messages["WARN"])) + assert.Equal(t, 1, len(baseLogger.messages["ERROR"])) // Verifies base logger is delegated to even when queue is unavailable. + assert.Equal(t, 0, len(baseLogger.messages["PANIC"])) + assert.Equal(t, "error:::::7a", baseLogger.messages["ERROR"][0]) + + assert.Equal(t, 0, len(store.storage[RecoverableErrorType])) + assert.Equal(t, 0, len(store.storage[PanicErrorType])) +} + +func TestErrorStreamLoggerErrorfPushFailEmitsErrorToBaseLogger(t *testing.T) { + baseLogger := NewMapStashLogger() + errMsg := "ERROR: this error should be emitted to the base logger." + store := NewMapStashErrorStore(true, errors.New(errMsg)) // Simulate error condition during push. + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Errorf("error:::::%d%s", 7, "aa") + + assert.Equal(t, 0, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 0, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 0, len(baseLogger.messages["INFO"])) + assert.Equal(t, 0, len(baseLogger.messages["WARN"])) + assert.Equal(t, 2, len(baseLogger.messages["ERROR"])) // Base logger should have a record of both errors and emit both. + assert.Equal(t, 0, len(baseLogger.messages["PANIC"])) + assert.Equal(t, "error:::::7aa", baseLogger.messages["ERROR"][0]) + assert.True(t, strings.Contains(baseLogger.messages["ERROR"][1], errMsg)) // Check original push error message wrapped. + assert.True(t, strings.Contains(baseLogger.messages["ERROR"][1], "Failed during push to store=")) // Check error contains info that it occurred during a push. + + assert.Equal(t, 0, len(store.storage[RecoverableErrorType])) + assert.Equal(t, 0, len(store.storage[PanicErrorType])) +} + +func TestErrorStreamLoggerPanicfWithAvailableQueue(t *testing.T) { + baseLogger := NewMapStashLogger() + store := NewMapStashErrorStore(true, nil) // Queue is available so a push actually occurs. + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Panicf("panic::::::%d", 8) + + assert.Equal(t, 0, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 0, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 0, len(baseLogger.messages["INFO"])) + assert.Equal(t, 0, len(baseLogger.messages["WARN"])) + assert.Equal(t, 0, len(baseLogger.messages["ERROR"])) + assert.Equal(t, 1, len(baseLogger.messages["PANIC"])) + assert.Equal(t, "panic::::::8", baseLogger.messages["PANIC"][0]) + + assert.Equal(t, 0, len(store.storage[RecoverableErrorType])) // Verifies a push occurred. + assert.Equal(t, 1, len(store.storage[PanicErrorType])) + assert.Equal(t, "panic::::::8", store.storage[PanicErrorType][0]) +} + +func TestErrorStreamLoggerPanicfNoAvailableQueue(t *testing.T) { + baseLogger := NewMapStashLogger() + store := NewMapStashErrorStore(false, nil) // Queue is unavailable. + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Panicf("panic::::::%d%s", 8, "b") + + assert.Equal(t, 0, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 0, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 0, len(baseLogger.messages["INFO"])) + assert.Equal(t, 0, len(baseLogger.messages["WARN"])) + assert.Equal(t, 0, len(baseLogger.messages["ERROR"])) + assert.Equal(t, 1, len(baseLogger.messages["PANIC"])) // Verifies base logger is delegated to even when queue is unavailable. + assert.Equal(t, "panic::::::8b", baseLogger.messages["PANIC"][0]) + + assert.Equal(t, 0, len(store.storage[RecoverableErrorType])) + assert.Equal(t, 0, len(store.storage[PanicErrorType])) +} + +func TestErrorStreamLoggerPanicfPushFailEmitsErrorToBaseLogger(t *testing.T) { + baseLogger := NewMapStashLogger() + errMsg := "PANIC: this error should be emitted to the base logger." + store := NewMapStashErrorStore(true, errors.New(errMsg)) // Simulate error condition during push. + + logger := NewErrorStreamLogger(baseLogger, store) + logger.Panicf("panic::::::%d%s", 8, "bb") + + assert.Equal(t, 0, len(baseLogger.messages["PRINTF"])) + assert.Equal(t, 0, len(baseLogger.messages["DEBUG"])) + assert.Equal(t, 0, len(baseLogger.messages["INFO"])) + assert.Equal(t, 0, len(baseLogger.messages["WARN"])) + assert.Equal(t, 1, len(baseLogger.messages["ERROR"])) // Base logger should have a record of the error due to a push. + assert.Equal(t, 1, len(baseLogger.messages["PANIC"])) // Base logger should have a record of the original panic/runtime error. + assert.Equal(t, "panic::::::8bb", baseLogger.messages["PANIC"][0]) + assert.True(t, strings.Contains(baseLogger.messages["ERROR"][0], errMsg)) // Check original push error message wrapped. + assert.True(t, strings.Contains(baseLogger.messages["ERROR"][0], "Failed during push to store=")) // Check error contains info that it occurred during a push. + + assert.Equal(t, 0, len(store.storage[RecoverableErrorType])) + assert.Equal(t, 0, len(store.storage[PanicErrorType])) +} diff --git a/idk/kinesis/source.go b/idk/kinesis/source.go index 4d5631b83..a52fc6fd1 100644 --- a/idk/kinesis/source.go +++ b/idk/kinesis/source.go @@ -15,6 +15,7 @@ import ( "github.com/aws/aws-sdk-go/service/kinesis/kinesisiface" "github.com/aws/aws-sdk-go/service/s3" "github.com/aws/aws-sdk-go/service/s3/s3iface" + "github.com/aws/aws-sdk-go/service/sqs" "github.com/molecula/featurebase/v3/idk" "github.com/molecula/featurebase/v3/idk/internal" @@ -37,6 +38,8 @@ type Source struct { StreamName string OffsetsPath string + ErrorQueueName string + schema []idk.Field paths idk.PathTable @@ -172,11 +175,17 @@ func (s *Source) initAWS() error { s.Log.Infof("Overriding default AWS region: %s", s.AWSRegion) config.Region = aws.String(s.AWSRegion) } + sess, err := session.NewSession(config) if err != nil { return errors.Wrap(err, "creating AWS session") } s.session = sess + + // Wrap the Logger instance to additionally broadcast errors and panics into an SQS queue. + queue := sqs.New(sess) + s.Log = NewErrorStreamLogger(s.Log, SinkErrorQueueFrom(queue, s)) + s.s3client = s3.New(sess) s.kinesisClient = kinesis.New(sess) return nil