Files
query-orchestration/internal/queryQueue/execute.go
T

92 lines
2.0 KiB
Go
Raw Normal View History

2025-01-07 16:30:45 +00:00
package queryqueue
2024-12-20 17:54:05 +00:00
import (
"context"
2024-12-24 17:13:48 +00:00
"fmt"
contextfull "queryorchestration/internal/contextFull"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
jsonextractor "queryorchestration/internal/jsonExtractor"
2025-01-03 13:41:07 +00:00
queryprocessor "queryorchestration/internal/queryProcessor"
2024-12-24 17:13:48 +00:00
"queryorchestration/internal/result"
2024-12-20 17:54:05 +00:00
"github.com/jackc/pgx/v5/pgtype"
)
func (q *Queue) Execute(ctx context.Context) error {
if q.unsyncedQueue == nil {
return nil
}
2025-01-06 15:31:39 +00:00
for _, query := range q.unsyncedQueue {
2024-12-20 17:54:05 +00:00
err := q.executeQuery(ctx, query)
if err != nil {
return err
}
}
q.unsyncedQueue = nil
return nil
}
2025-01-06 15:31:39 +00:00
func (q *Queue) executeQuery(ctx context.Context, qu *queryprocessor.Query) error {
2025-01-03 13:41:07 +00:00
resultIDs := make([]pgtype.UUID, len(qu.RequiredQueryIDs))
for index, id := range qu.RequiredQueryIDs {
2024-12-20 17:54:05 +00:00
var queryVersion int32
2025-01-06 15:31:39 +00:00
for _, entry := range q.collectorQueries {
2024-12-20 17:54:05 +00:00
if entry.ID == id {
queryVersion = entry.Version
break
}
}
2025-01-06 15:31:39 +00:00
for _, entry := range q.results {
2024-12-20 17:54:05 +00:00
if entry.QueryID == id && entry.QueryVersion == queryVersion {
resultIDs[index] = database.MustToDBUUID(entry.ID)
break
}
}
}
2025-01-06 12:26:28 +00:00
values, err := q.db.Queries.ListResultValuesByID(ctx, resultIDs)
2024-12-20 17:54:05 +00:00
if err != nil {
return err
}
2024-12-23 12:26:05 +00:00
cleanValues := make([]result.Value, len(values))
2024-12-24 17:13:48 +00:00
for index, r := range values {
cleanValue, err := q.getResultValue(&r)
if err != nil {
return err
}
cleanValues[index] = cleanValue
2024-12-20 17:54:05 +00:00
}
2025-01-06 15:31:39 +00:00
err = q.setResult(ctx, qu, cleanValues)
2024-12-20 17:54:05 +00:00
if err != nil {
return err
}
return nil
}
2024-12-24 17:13:48 +00:00
func (q *Queue) getResultValue(res *repository.ListResultValuesByIDRow) (result.Value, error) {
2025-01-03 13:41:07 +00:00
var queryType queryprocessor.Type
2025-01-06 15:31:39 +00:00
for _, qu := range q.collectorQueries {
2024-12-24 17:13:48 +00:00
if qu.ID == database.MustToUUID(res.Queryid) {
queryType = qu.Type
}
}
switch queryType {
2025-01-03 13:41:07 +00:00
case queryprocessor.TypeJsonExtractor:
2024-12-24 17:13:48 +00:00
return jsonextractor.NewResult(res.Value), nil
2025-01-03 13:41:07 +00:00
case queryprocessor.TypeContextFull:
2024-12-24 17:13:48 +00:00
return contextfull.NewResult(res.Value), nil
default:
return nil, fmt.Errorf("attempting to process invalid query type")
}
}