Merged in feature/splitprocess (pull request #83)
Split Queue Process * queue
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user