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
|
|
|
"queryorchestration/internal/database"
|
2025-01-17 12:00:32 +00:00
|
|
|
queryprocessor "queryorchestration/internal/query/processor"
|
2024-12-20 17:54:05 +00:00
|
|
|
)
|
|
|
|
|
|
|
|
|
|
func (q *Queue) getUnsyncedQueries() {
|
2025-01-06 15:31:39 +00:00
|
|
|
for _, query := range q.collectorQueries {
|
2025-01-08 18:14:15 +00:00
|
|
|
if q.isQuerySynced(query) {
|
2024-12-20 17:54:05 +00:00
|
|
|
continue
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-06 15:31:39 +00:00
|
|
|
q.Add(query)
|
2024-12-20 17:54:05 +00:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-08 18:14:15 +00:00
|
|
|
func (q *Queue) isQuerySynced(query *queryprocessor.Query) bool {
|
|
|
|
|
for _, result := range q.results {
|
|
|
|
|
if result.QueryID == query.ID && result.QueryVersion == query.Version {
|
|
|
|
|
return true
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return false
|
|
|
|
|
}
|
|
|
|
|
|
2024-12-20 17:54:05 +00:00
|
|
|
func (c *Queue) getCollectorQueries(ctx context.Context) error {
|
|
|
|
|
if c.collectorQueries != nil {
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
2024-12-23 18:13:57 +00:00
|
|
|
id := database.MustToDBUUID(c.collector.ID)
|
2024-12-20 17:54:05 +00:00
|
|
|
|
2025-01-06 12:26:28 +00:00
|
|
|
queries, err := c.db.Queries.GetCollectorQueries(ctx, id)
|
2024-12-20 17:54:05 +00:00
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-06 15:31:39 +00:00
|
|
|
cleanQueries := make([]*queryprocessor.Query, len(queries))
|
2024-12-20 17:54:05 +00:00
|
|
|
for index, dbQuery := range queries {
|
2025-01-15 12:19:49 +00:00
|
|
|
cleanQuery, err := queryprocessor.ParseDBCollectorQuery(dbQuery)
|
2024-12-20 17:54:05 +00:00
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-06 15:31:39 +00:00
|
|
|
cleanQueries[index] = cleanQuery
|
2024-12-20 17:54:05 +00:00
|
|
|
}
|
|
|
|
|
|
2025-01-06 15:31:39 +00:00
|
|
|
c.collectorQueries = cleanQueries
|
2024-12-20 17:54:05 +00:00
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-03 13:41:07 +00:00
|
|
|
func (q *Queue) Add(qu *queryprocessor.Query) {
|
2025-01-06 15:31:39 +00:00
|
|
|
dependentQueries := []*queryprocessor.Query{}
|
2024-12-20 17:54:05 +00:00
|
|
|
requiredIndex := -1
|
|
|
|
|
|
|
|
|
|
if q.unsyncedQueue == nil {
|
2025-01-06 15:31:39 +00:00
|
|
|
q.unsyncedQueue = []*queryprocessor.Query{}
|
2024-12-20 17:54:05 +00:00
|
|
|
} else {
|
2025-01-06 15:31:39 +00:00
|
|
|
for index, entry := range q.unsyncedQueue {
|
2024-12-24 12:47:40 +00:00
|
|
|
if entry.ID == qu.ID {
|
2024-12-20 17:54:05 +00:00
|
|
|
return
|
|
|
|
|
}
|
2025-01-20 13:31:48 +00:00
|
|
|
if qu.RequiredQueryIDs != nil {
|
|
|
|
|
for _, id := range *qu.RequiredQueryIDs {
|
|
|
|
|
if entry.ID == id {
|
|
|
|
|
requiredIndex = index
|
|
|
|
|
break
|
|
|
|
|
}
|
2025-01-03 13:41:07 +00:00
|
|
|
}
|
2024-12-20 17:54:05 +00:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2025-01-06 15:31:39 +00:00
|
|
|
for _, entry := range q.collectorQueries {
|
2025-01-20 13:31:48 +00:00
|
|
|
if entry.RequiredQueryIDs != nil {
|
|
|
|
|
for _, id := range *entry.RequiredQueryIDs {
|
|
|
|
|
if qu.ID == id {
|
|
|
|
|
dependentQueries = append(dependentQueries, entry)
|
|
|
|
|
break
|
|
|
|
|
}
|
2025-01-03 13:41:07 +00:00
|
|
|
}
|
2024-12-20 17:54:05 +00:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if requiredIndex != -1 {
|
2025-01-06 15:31:39 +00:00
|
|
|
q.unsyncedQueue = append((q.unsyncedQueue)[:requiredIndex+1], append([]*queryprocessor.Query{qu}, (q.unsyncedQueue)[requiredIndex+1:]...)...)
|
2024-12-20 17:54:05 +00:00
|
|
|
} else {
|
2025-01-06 15:31:39 +00:00
|
|
|
q.unsyncedQueue = append([]*queryprocessor.Query{qu}, q.unsyncedQueue...)
|
2024-12-20 17:54:05 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
for _, entry := range dependentQueries {
|
2025-01-06 15:31:39 +00:00
|
|
|
q.Add(entry)
|
2024-12-20 17:54:05 +00:00
|
|
|
}
|
|
|
|
|
}
|