From 4b8c930ff1b99e2af745a32fec8b26538906e6a3 Mon Sep 17 00:00:00 2001 From: Michael McGuinness Date: Tue, 4 Mar 2025 12:14:32 +0000 Subject: [PATCH] Merged in feature/splitprocess (pull request #83) Split Queue Process * queue --- internal/server/runner/poll.go | 33 ++++++---- internal/server/runner/poll_test.go | 94 ++++++++++++++++++++++++++--- 2 files changed, 107 insertions(+), 20 deletions(-) diff --git a/internal/server/runner/poll.go b/internal/server/runner/poll.go index b7f30916..a3f0ab36 100644 --- a/internal/server/runner/poll.go +++ b/internal/server/runner/poll.go @@ -5,6 +5,8 @@ import ( "fmt" "log/slog" "queryorchestration/internal/serviceconfig/queue" + + "github.com/aws/aws-sdk-go-v2/service/sqs/types" ) type Server struct { @@ -43,20 +45,29 @@ func (c *Server) pollMessage(ctx context.Context) error { } for _, message := range result.Messages { - slog.Debug("processing message", "id", *message.MessageId, "body", *message.Body) - - err := c.cfg.GetController().Process(ctx, &message) + err := c.processMessage(ctx, &message) if err != nil { slog.Error("message process fail", "id", *message.MessageId, "err", err) } - - err = c.cfg.DeleteFromQueue(ctx, &queue.DeleteParams{ - QueueURL: c.cfg.GetQueueURL(), - ReceiptHandle: message.ReceiptHandle, - }) - if err != nil { - slog.Error("message delete fail", "id", *message.MessageId, "err", err) - } + } + + return nil +} + +func (c *Server) processMessage(ctx context.Context, message *types.Message) error { + slog.Debug("processing message", "id", *message.MessageId, "body", *message.Body) + + err := c.cfg.GetController().Process(ctx, message) + if err != nil { + return fmt.Errorf("message process fail: %v", err) + } + + err = c.cfg.DeleteFromQueue(ctx, &queue.DeleteParams{ + QueueURL: c.cfg.GetQueueURL(), + ReceiptHandle: message.ReceiptHandle, + }) + if err != nil { + return fmt.Errorf("message delete fail: %v", err) } return nil diff --git a/internal/server/runner/poll_test.go b/internal/server/runner/poll_test.go index b5b243fe..379072d7 100644 --- a/internal/server/runner/poll_test.go +++ b/internal/server/runner/poll_test.go @@ -16,12 +16,11 @@ import ( func TestPollMessages(t *testing.T) { ctx := context.Background() - mockSQS := queuemock.NewMockSQSClient(t) cfg := &BaseConfig{} + mockSQS := queuemock.NewMockSQSClient(t) cfg.QueueClient = mockSQS - cfg.ControllerFunc = func() Controller { - return runnermock.NewMockController(t) - } + mockController := runnermock.NewMockController(t) + cfg.Controller = mockController cfg.QueueURL = "/i/am/here" res := sqs.ReceiveMessageOutput{Messages: []types.Message{}} @@ -49,12 +48,11 @@ func TestPollMessages(t *testing.T) { func TestPollMessage(t *testing.T) { ctx := context.Background() - mockSQS := queuemock.NewMockSQSClient(t) cfg := &BaseConfig{} + mockSQS := queuemock.NewMockSQSClient(t) cfg.QueueClient = mockSQS - cfg.ControllerFunc = func() Controller { - return runnermock.NewMockController(t) - } + mockController := runnermock.NewMockController(t) + cfg.Controller = mockController cfg.QueueURL = "/i/am/here" ser := &Server{ @@ -62,6 +60,10 @@ func TestPollMessage(t *testing.T) { cleanup: func() error { return nil }, } + id := "example_id" + body := "example_body" + receipt := "receipts" + mockSQS.EXPECT(). ReceiveMessage( mock.Anything, @@ -70,8 +72,82 @@ func TestPollMessage(t *testing.T) { }), mock.Anything, ). - Return(&sqs.ReceiveMessageOutput{}, nil) + Return(&sqs.ReceiveMessageOutput{ + Messages: []types.Message{ + { + MessageId: &id, + Body: &body, + ReceiptHandle: &receipt, + }, + }, + }, nil) + + mockController.EXPECT(). + Process( + mock.Anything, + mock.MatchedBy(func(in *types.Message) bool { + return *in.MessageId == id && *in.Body == body + }), + ). + Return(nil) + + mockSQS.EXPECT(). + DeleteMessage( + mock.Anything, + mock.MatchedBy(func(in *sqs.DeleteMessageInput) bool { + return *in.QueueUrl == cfg.QueueURL && *in.ReceiptHandle == receipt + }), + mock.Anything, + ). + Return(&sqs.DeleteMessageOutput{}, nil) err := ser.pollMessage(ctx) assert.NoError(t, err) } + +func TestProcessMessage(t *testing.T) { + ctx := context.Background() + + cfg := &BaseConfig{} + mockSQS := queuemock.NewMockSQSClient(t) + cfg.QueueClient = mockSQS + mockController := runnermock.NewMockController(t) + cfg.Controller = mockController + cfg.QueueURL = "/i/am/here" + + ser := &Server{ + cfg: cfg, + cleanup: func() error { return nil }, + } + + id := "example_id" + body := "example_body" + receipt := "receipts" + msg := &types.Message{ + MessageId: &id, + Body: &body, + ReceiptHandle: &receipt, + } + + mockController.EXPECT(). + Process( + mock.Anything, + mock.MatchedBy(func(in *types.Message) bool { + return *in.MessageId == id && *in.Body == body + }), + ). + Return(nil) + + mockSQS.EXPECT(). + DeleteMessage( + mock.Anything, + mock.MatchedBy(func(in *sqs.DeleteMessageInput) bool { + return *in.QueueUrl == cfg.QueueURL && *in.ReceiptHandle == receipt + }), + mock.Anything, + ). + Return(&sqs.DeleteMessageOutput{}, nil) + + err := ser.processMessage(ctx, msg) + assert.NoError(t, err) +}