package queryqueue import ( "context" "queryorchestration/internal/database" queryprocessor "queryorchestration/internal/query/processor" ) func (q *Queue) getUnsyncedQueries() { for _, query := range q.collectorQueries { if q.isQuerySynced(query) { continue } q.Add(query) } } 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 } func (c *Queue) getCollectorQueries(ctx context.Context) error { if c.collectorQueries != nil { return nil } id := database.MustToDBUUID(c.collector.ID) queries, err := c.db.Queries.ListCollectorQueries(ctx, id) if err != nil { return err } cleanQueries := make([]*queryprocessor.Query, len(queries)) for index, dbQuery := range queries { cleanQuery, err := queryprocessor.ParseDBCollectorQuery(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 } if qu.RequiredQueryIDs != nil { for _, id := range *qu.RequiredQueryIDs { if entry.ID == id { requiredIndex = index break } } } } } for _, entry := range q.collectorQueries { if entry.RequiredQueryIDs != nil { 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) } }