78dc2f6cbc
Add Localstack to Compose * localstack * addedlocalstackpluscleanup#
64 lines
1.2 KiB
Go
64 lines
1.2 KiB
Go
package runner
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log/slog"
|
|
"queryorchestration/internal/serviceconfig/queue"
|
|
)
|
|
|
|
type Server struct {
|
|
cleanup func() error
|
|
cfg ListenerConfig
|
|
}
|
|
|
|
func (c *Server) 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) 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: %v", err)
|
|
}
|
|
|
|
for _, message := range result.Messages {
|
|
slog.Info("processing message", "id", message.MessageId)
|
|
|
|
err := c.cfg.GetController().Process(ctx, &message)
|
|
if err != nil {
|
|
slog.Error("message process fail", "err", err)
|
|
}
|
|
|
|
err = c.cfg.DeleteFromQueue(ctx, &queue.DeleteParams{
|
|
QueueURL: c.cfg.GetQueueURL(),
|
|
ReceiptHandle: message.ReceiptHandle,
|
|
})
|
|
if err != nil {
|
|
slog.Error("message delete fail", "err", err)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|