Files
query-orchestration/internal/server/runner/poll.go
T
Michael McGuinness 71f9802e1a Merged in feature/splitqueryrunning (pull request #57)
Split Query Running + Debugging Full Flow

* completedquerysyncrunner

* spliitinglogic

* synccomplete

* informdependents

* only push same collector

* deps

* livetesting

* foundissue

* some issues resolved

* activeupdate

* collectorupdatefixes

* fix dbquesries

* tests

* tests

* pollingdebug
2025-02-11 15:22:59 +00:00

64 lines
1.3 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.Debug("processing message", "id", *message.MessageId, "body", *message.Body)
err := c.cfg.GetController().Process(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
}