Files
query-orchestration/internal/server/runner/poll.go
T
Michael McGuinness 3028fe7eaa Merged in feature/linting (pull request #168)
Linting Updates

* precommit

* smallerchecks

* govuln
2025-06-23 15:58:20 +00:00

110 lines
2.3 KiB
Go

package runner
import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"queryorchestration/internal/serviceconfig/queue"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)
type Server[B any] struct {
cleanup func() error
cfg ConfigProvider[B]
}
func (c *Server[B]) Listen(ctx context.Context) {
slog.Info("Listening to queue")
defer func() {
if err := c.cleanup(); err != nil {
fmt.Printf("Error during cleanup: %v", err)
}
}()
for {
select {
case <-ctx.Done():
return
default:
err := c.pollMessage(ctx)
if err != nil {
slog.Error(err.Error())
}
}
}
}
func (c *Server[B]) pollMessage(ctx context.Context) error {
result, err := c.cfg.ReceiveFromQueue(ctx, &queue.ReceiveParams{
QueueURL: c.cfg.GetQueueURL(),
})
if err != nil {
return fmt.Errorf("message fetch fail: %w", err)
}
for _, message := range result.Messages {
err := c.processMessage(ctx, &message)
if err != nil {
slog.Error("message process fail", "id", *message.MessageId, "err", err)
}
}
return nil
}
func (c *Server[B]) processMessage(ctx context.Context, message *types.Message) error {
body := c.processBody(message)
if body == nil {
return nil
}
slog.Debug("processing message", "id", *message.MessageId, "body", *message.Body)
isProcessed := c.cfg.GetController().Process(ctx, *body)
if !isProcessed {
return errors.New("message process fail")
}
err := c.cfg.DeleteFromQueue(ctx, &queue.DeleteParams{
QueueURL: c.cfg.GetQueueURL(),
ReceiptHandle: message.ReceiptHandle,
})
if err != nil {
return fmt.Errorf("message delete fail: %w", err)
}
return nil
}
func (c *Server[B]) processBody(message *types.Message) *B {
if message == nil {
slog.Error("nil message")
return nil
} else if message.MessageId == nil {
slog.Error("nil message id")
return nil
} else if message.Body == nil {
slog.Error("nil body", "messageId", *message.MessageId)
return nil
}
var body B
err := json.Unmarshal([]byte(*message.Body), &body)
if err != nil {
slog.Error("parsing body", "messageId", *message.MessageId, "body", *message.Body)
return nil
}
err = c.cfg.GetValidator().Struct(body)
if err != nil {
slog.Error("invalid body", "messageId", *message.MessageId, "error", err)
return nil
}
return &body
}