Files
query-orchestration/internal/query/queue/execute.go
T
Michael McGuinness fa95d733ca Merged in feature/shorttestsanddirtidy (pull request #26)
Add short tests and Tidy internal directories

* complete the tasks
2025-01-17 12:00:32 +00:00

70 lines
1.3 KiB
Go

package queryqueue
import (
"context"
"queryorchestration/internal/database"
queryprocessor "queryorchestration/internal/query/processor"
"queryorchestration/internal/query/result"
"github.com/jackc/pgx/v5/pgtype"
)
func (q *Queue) Execute(ctx context.Context) error {
if q.unsyncedQueue == nil {
return nil
}
for _, query := range q.unsyncedQueue {
err := q.executeQuery(ctx, query)
if err != nil {
return err
}
}
q.unsyncedQueue = nil
return nil
}
func (q *Queue) executeQuery(ctx context.Context, qu *queryprocessor.Query) error {
resultIDs := make([]pgtype.UUID, len(qu.RequiredQueryIDs))
for index, id := range qu.RequiredQueryIDs {
var queryVersion int32
for _, entry := range q.collectorQueries {
if entry.ID == id {
queryVersion = entry.Version
break
}
}
for _, entry := range q.results {
if entry.QueryID == id && entry.QueryVersion == queryVersion {
resultIDs[index] = database.MustToDBUUID(entry.ID)
break
}
}
}
values, err := q.db.Queries.ListResultValuesByID(ctx, resultIDs)
if err != nil {
return err
}
cleanValues := make([]result.Value, len(values))
for index, r := range values {
cleanValue, err := q.getResultValue(r)
if err != nil {
return err
}
cleanValues[index] = cleanValue
}
err = q.setResult(ctx, qu, cleanValues)
if err != nil {
return err
}
return nil
}