Files
query-orchestration/internal/query/create.go
T
2024-12-20 17:54:05 +00:00

89 lines
1.7 KiB
Go

package query
import (
"context"
"gotemplate/internal/database"
)
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, err := database.ToDBUUID(c.collector.ID)
if err != nil {
return err
}
queries, err := c.db.GetCollectorQueries(ctx, id)
if err != nil {
return err
}
cleanQueries := make([]Query, len(queries))
for index, dbQuery := range queries {
cleanQuery, err := ParseDBQuery(&dbQuery)
if err != nil {
return err
}
cleanQueries[index] = *cleanQuery
}
c.collectorQueries = &cleanQueries
return nil
}
func (q *Queue) Add(query *Query) {
dependentQueries := []Query{}
requiredIndex := -1
if q.unsyncedQueue == nil {
q.unsyncedQueue = &[]Query{}
} else {
for index, entry := range *q.unsyncedQueue {
if entry.ID == query.ID {
return
}
if entry.ID == query.RequiredQueryID {
requiredIndex = index
}
}
}
for _, entry := range *q.collectorQueries {
if entry.RequiredQueryID == query.ID {
dependentQueries = append(dependentQueries, entry)
}
}
if requiredIndex != -1 {
*q.unsyncedQueue = append((*q.unsyncedQueue)[:requiredIndex], append([]Query{*query}, (*q.unsyncedQueue)[requiredIndex:]...)...)
} else {
*q.unsyncedQueue = append([]Query{*query}, *q.unsyncedQueue...)
}
for _, entry := range dependentQueries {
q.Add(&entry)
}
}