package queryQueue import ( "context" "queryorchestration/internal/database" "queryorchestration/internal/query" ) 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([]query.QueryRow, len(queries)) for index, dbQuery := range queries { cleanQuery, err := query.ParseDBQueryRow(&dbQuery) if err != nil { return err } cleanQueries[index] = *cleanQuery } c.collectorQueries = &cleanQueries return nil } func (q *Queue) Add(qu *query.QueryRow) { dependentQueries := []query.QueryRow{} requiredIndex := -1 if q.unsyncedQueue == nil { q.unsyncedQueue = &[]query.QueryRow{} } else { for index, entry := range *q.unsyncedQueue { if entry.ID == qu.ID { return } if entry.ID == qu.RequiredQueryID { requiredIndex = index } } } for _, entry := range *q.collectorQueries { if entry.RequiredQueryID == qu.ID { dependentQueries = append(dependentQueries, entry) } } if requiredIndex != -1 { *q.unsyncedQueue = append((*q.unsyncedQueue)[:requiredIndex+1], append([]query.QueryRow{*qu}, (*q.unsyncedQueue)[requiredIndex+1:]...)...) } else { *q.unsyncedQueue = append([]query.QueryRow{*qu}, *q.unsyncedQueue...) } for _, entry := range dependentQueries { q.Add(&entry) } }