2024-12-24 12:47:40 +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
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
for _, query := range *q.unsyncedQueue {
|
|
|
|
|
err := q.executeQuery(ctx, query)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
q.unsyncedQueue = nil
|
|
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-03 13:41:07 +00:00
|
|
|
func (q *Queue) executeQuery(ctx context.Context, qu queryprocessor.Query) error {
|
|
|
|
|
resultIDs := make([]pgtype.UUID, len(qu.RequiredQueryIDs))
|
|
|
|
|
for index, id := range qu.RequiredQueryIDs {
|
2024-12-20 17:54:05 +00:00
|
|
|
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
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
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
|
|
|
}
|
|
|
|
|
|
2024-12-24 12:47:40 +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
|
2024-12-24 17:13:48 +00:00
|
|
|
for _, qu := range *q.collectorQueries {
|
|
|
|
|
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")
|
|
|
|
|
}
|
|
|
|
|
}
|