93 lines
1.9 KiB
Go
93 lines
1.9 KiB
Go
package queryQueue
|
|
|
|
import (
|
|
"context"
|
|
"queryorchestration/internal/database"
|
|
queryprocessor "queryorchestration/internal/queryProcessor"
|
|
)
|
|
|
|
func (q *Queue) getUnsyncedQueries() {
|
|
for _, query := range *q.collectorQueries {
|
|
isSynced := false
|
|
for _, result := range *q.results {
|
|
if result.QueryID != query.ID || result.QueryVersion != query.Version {
|
|
continue
|
|
}
|
|
isSynced = true
|
|
break
|
|
}
|
|
if isSynced {
|
|
continue
|
|
}
|
|
|
|
q.Add(&query)
|
|
}
|
|
}
|
|
|
|
func (c *Queue) getCollectorQueries(ctx context.Context) error {
|
|
if c.collectorQueries != nil {
|
|
return nil
|
|
}
|
|
|
|
id := database.MustToDBUUID(c.collector.ID)
|
|
|
|
queries, err := c.db.GetCollectorQueries(ctx, id)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
cleanQueries := make([]queryprocessor.Query, len(queries))
|
|
for index, dbQuery := range queries {
|
|
cleanQuery, err := queryprocessor.ParseDBQuery(&dbQuery)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
cleanQueries[index] = *cleanQuery
|
|
}
|
|
|
|
c.collectorQueries = &cleanQueries
|
|
|
|
return nil
|
|
}
|
|
|
|
func (q *Queue) Add(qu *queryprocessor.Query) {
|
|
dependentQueries := []queryprocessor.Query{}
|
|
requiredIndex := -1
|
|
|
|
if q.unsyncedQueue == nil {
|
|
q.unsyncedQueue = &[]queryprocessor.Query{}
|
|
} else {
|
|
for index, entry := range *q.unsyncedQueue {
|
|
if entry.ID == qu.ID {
|
|
return
|
|
}
|
|
for _, id := range qu.RequiredQueryIDs {
|
|
if entry.ID == id {
|
|
requiredIndex = index
|
|
break
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
for _, entry := range *q.collectorQueries {
|
|
for _, id := range entry.RequiredQueryIDs {
|
|
if qu.ID == id {
|
|
dependentQueries = append(dependentQueries, entry)
|
|
break
|
|
}
|
|
}
|
|
}
|
|
|
|
if requiredIndex != -1 {
|
|
*q.unsyncedQueue = append((*q.unsyncedQueue)[:requiredIndex+1], append([]queryprocessor.Query{*qu}, (*q.unsyncedQueue)[requiredIndex+1:]...)...)
|
|
} else {
|
|
*q.unsyncedQueue = append([]queryprocessor.Query{*qu}, *q.unsyncedQueue...)
|
|
}
|
|
|
|
for _, entry := range dependentQueries {
|
|
q.Add(&entry)
|
|
}
|
|
}
|