From b888e3450fefa6519d7bb3f422ec65d63ac659a6 Mon Sep 17 00:00:00 2001 From: Michael McGuinness Date: Mon, 20 Jan 2025 13:51:22 +0000 Subject: [PATCH] Merged in feature/outputqueue (pull request #29) Remove Queue Controller Config * removeout --- api/queryRunner/queryrunner.go | 20 ++------------------ api/queryRunner/queryrunner_test.go | 11 +---------- internal/server/queue/config.go | 2 +- internal/server/queue/poll.go | 2 +- internal/server/queue/poll_test.go | 2 +- 5 files changed, 6 insertions(+), 31 deletions(-) diff --git a/api/queryRunner/queryrunner.go b/api/queryRunner/queryrunner.go index 3b5a2acd..9992564a 100644 --- a/api/queryRunner/queryrunner.go +++ b/api/queryRunner/queryrunner.go @@ -4,11 +4,9 @@ import ( "context" "encoding/json" "queryorchestration/internal/query/document" - "queryorchestration/internal/server/queue" "github.com/go-playground/validator/v10" - "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/service/sqs/types" "github.com/google/uuid" ) @@ -29,9 +27,9 @@ type DocumentQueryEvent struct { ID uuid.UUID `json:"id"` } -func (s QueryRunner) Process(ctx context.Context, config *queue.Config, msg *types.Message) error { +func (s QueryRunner) Process(ctx context.Context, req *types.Message) error { var body document.Document - err := json.Unmarshal([]byte(*msg.Body), &body) + err := json.Unmarshal([]byte(*req.Body), &body) if err != nil { return err } @@ -46,19 +44,5 @@ func (s QueryRunner) Process(ctx context.Context, config *queue.Config, msg *typ return err } - queryEvent := DocumentQueryEvent{ - ID: body.ID, - } - - err = queue.Send(ctx, config, queryEvent, map[string]types.MessageAttributeValue{ - "type": { - DataType: aws.String("String"), - StringValue: aws.String("DOCQUERY"), - }, - }) - if err != nil { - return err - } - return nil } diff --git a/api/queryRunner/queryrunner_test.go b/api/queryRunner/queryrunner_test.go index 8e8bdca0..67c489cf 100644 --- a/api/queryRunner/queryrunner_test.go +++ b/api/queryRunner/queryrunner_test.go @@ -7,8 +7,6 @@ import ( "queryorchestration/internal/database" "queryorchestration/internal/database/repository" "queryorchestration/internal/query/document" - "queryorchestration/internal/server/queue" - "queryorchestration/internal/test" "testing" "github.com/aws/aws-sdk-go-v2/service/sqs/types" @@ -25,9 +23,6 @@ func TestQueryRunner(t *testing.T) { } ctx := context.Background() - qCfg, cleanup := test.CreateQueue(t, ctx, &test.CreateQueueConfig{}) - defer cleanup() - pool, err := pgxmock.NewPool() if err != nil { t.Fatalf("failed to open pgxmock database: %v", err) @@ -61,10 +56,6 @@ func TestQueryRunner(t *testing.T) { minCleanVersion := int32(1) minTextVersion := int32(1) - cfg := &queue.Config{ - URL: qCfg.URL, - Client: qCfg.Client, - } qV := int32(1) pool.ExpectQuery("name: GetCollectorFromJobID :one").WithArgs(database.MustToDBUUID(doc.JobID)). @@ -83,6 +74,6 @@ func TestQueryRunner(t *testing.T) { AddRow(collectorID, queryID, repository.NullQuerytype{Querytype: repository.QuerytypeContextFull, Valid: true}, &qV, []pgtype.UUID{}), ) - err = runner.Process(ctx, cfg, msg) + err = runner.Process(ctx, msg) assert.Nil(t, err) } diff --git a/internal/server/queue/config.go b/internal/server/queue/config.go index 1bfa85cb..9b74736e 100644 --- a/internal/server/queue/config.go +++ b/internal/server/queue/config.go @@ -13,5 +13,5 @@ type Config struct { } type Controller interface { - Process(ctx context.Context, config *Config, msg *types.Message) error + Process(ctx context.Context, message *types.Message) error } diff --git a/internal/server/queue/poll.go b/internal/server/queue/poll.go index 4e6c2d5c..4bff112d 100644 --- a/internal/server/queue/poll.go +++ b/internal/server/queue/poll.go @@ -35,7 +35,7 @@ func PollMessage(ctx context.Context, queueConfig *PollConfig) error { for _, message := range result.Messages { go func() { - err := queueConfig.Controller.Process(ctx, queueConfig.Config, &message) + err := queueConfig.Controller.Process(ctx, &message) if err != nil { log.Printf("message process fail: %v", err) } diff --git a/internal/server/queue/poll_test.go b/internal/server/queue/poll_test.go index dfcfeed6..c9616f33 100644 --- a/internal/server/queue/poll_test.go +++ b/internal/server/queue/poll_test.go @@ -13,7 +13,7 @@ import ( type MockController struct{} -func (s MockController) Process(ctx context.Context, config *queue.Config, msg *types.Message) error { +func (s MockController) Process(ctx context.Context, req *types.Message) error { return nil }