2024-12-24 12:47:40 +00:00
|
|
|
package queryQueue
|
2024-12-20 17:54:05 +00:00
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"context"
|
|
|
|
|
"gotemplate/internal/database"
|
2024-12-24 12:47:40 +00:00
|
|
|
"gotemplate/internal/query"
|
2024-12-20 17:54:05 +00:00
|
|
|
)
|
|
|
|
|
|
|
|
|
|
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
|
|
|
|
|
}
|
|
|
|
|
|
2024-12-23 18:13:57 +00:00
|
|
|
id := database.MustToDBUUID(c.collector.ID)
|
2024-12-20 17:54:05 +00:00
|
|
|
|
|
|
|
|
queries, err := c.db.GetCollectorQueries(ctx, id)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
2024-12-24 12:47:40 +00:00
|
|
|
cleanQueries := make([]query.Query, len(queries))
|
2024-12-20 17:54:05 +00:00
|
|
|
for index, dbQuery := range queries {
|
2024-12-24 12:47:40 +00:00
|
|
|
cleanQuery, err := query.ParseDBQuery(&dbQuery)
|
2024-12-20 17:54:05 +00:00
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
cleanQueries[index] = *cleanQuery
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
c.collectorQueries = &cleanQueries
|
|
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
2024-12-24 12:47:40 +00:00
|
|
|
func (q *Queue) Add(qu *query.Query) {
|
|
|
|
|
dependentQueries := []query.Query{}
|
2024-12-20 17:54:05 +00:00
|
|
|
requiredIndex := -1
|
|
|
|
|
|
|
|
|
|
if q.unsyncedQueue == nil {
|
2024-12-24 12:47:40 +00:00
|
|
|
q.unsyncedQueue = &[]query.Query{}
|
2024-12-20 17:54:05 +00:00
|
|
|
} else {
|
|
|
|
|
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
|
|
|
|
|
}
|
2024-12-24 12:47:40 +00:00
|
|
|
if entry.ID == qu.RequiredQueryID {
|
2024-12-20 17:54:05 +00:00
|
|
|
requiredIndex = index
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
for _, entry := range *q.collectorQueries {
|
2024-12-24 12:47:40 +00:00
|
|
|
if entry.RequiredQueryID == qu.ID {
|
2024-12-20 17:54:05 +00:00
|
|
|
dependentQueries = append(dependentQueries, entry)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if requiredIndex != -1 {
|
2024-12-24 12:47:40 +00:00
|
|
|
*q.unsyncedQueue = append((*q.unsyncedQueue)[:requiredIndex], append([]query.Query{*qu}, (*q.unsyncedQueue)[requiredIndex:]...)...)
|
2024-12-20 17:54:05 +00:00
|
|
|
} else {
|
2024-12-24 12:47:40 +00:00
|
|
|
*q.unsyncedQueue = append([]query.Query{*qu}, *q.unsyncedQueue...)
|
2024-12-20 17:54:05 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
for _, entry := range dependentQueries {
|
|
|
|
|
q.Add(&entry)
|
|
|
|
|
}
|
|
|
|
|
}
|