From f95114e11add3dfb2c6ada062df7fb0378e7def9 Mon Sep 17 00:00:00 2001 From: Michael McGuinness Date: Fri, 20 Dec 2024 17:35:33 +0000 Subject: [PATCH] queuecreation --- api/queue/document.go | 4 +- database/queries/result.sql | 8 +- internal/collector/service.go | 47 ++++ internal/database/repository/models.go | 2 +- internal/database/repository/result.sql.go | 69 +++++- internal/database/uuid.go | 15 ++ internal/document/service.go | 104 +++------ internal/query/database.go | 38 ++++ internal/query/service.go | 251 +++++++++++++++++++++ test/unit/internal/document/sync_test.go | 29 ++- test/unit/internal/query/queue_test.go | 125 ++++++++++ 11 files changed, 603 insertions(+), 89 deletions(-) create mode 100644 internal/collector/service.go create mode 100644 internal/query/database.go create mode 100644 internal/query/service.go create mode 100644 test/unit/internal/query/queue_test.go diff --git a/api/queue/document.go b/api/queue/document.go index b4e2092d..236d46a6 100644 --- a/api/queue/document.go +++ b/api/queue/document.go @@ -31,7 +31,9 @@ func (s *DocumentController) Sync(ctx context.Context, config *queue.QueueConfig return err } - err = s.document.Sync(ctx, body) + // TODO normalise input here + + err = s.document.Sync(ctx, &body) if err != nil { return err } diff --git a/database/queries/result.sql b/database/queries/result.sql index a564ad22..07c67bf1 100644 --- a/database/queries/result.sql +++ b/database/queries/result.sql @@ -1,2 +1,8 @@ -- name: ListResultsByDocumentID :many -SELECT id, queryId, cleanVersion, textVersion, queryVersion FROM results where documentId = $1 and cleanVersion >= $2 and textVersion >= $3; \ No newline at end of file +SELECT id, queryId, queryVersion FROM results where documentId = $1 and cleanVersion >= $2 and textVersion >= $3; + +-- name: ListResultValuesByID :many +SELECT id, queryId, value FROM results where id = ANY($1); + +-- name: SetResult :exec +INSERT INTO results (id, queryId, documentId, value, cleanVersion, textVersion, queryVersion) VALUES ($1, $2, $3, $4, $5, $6, $7); \ No newline at end of file diff --git a/internal/collector/service.go b/internal/collector/service.go new file mode 100644 index 00000000..36469d6e --- /dev/null +++ b/internal/collector/service.go @@ -0,0 +1,47 @@ +package collector + +import ( + "context" + "gotemplate/internal/database" + "gotemplate/internal/database/repository" + + "github.com/google/uuid" +) + +type Collector struct { + ID uuid.UUID + MinCleanVersion int32 + MinTextVersion int32 + db *repository.Queries +} + +func NewByJobId(ctx context.Context, db *repository.Queries, jobID uuid.UUID) (*Collector, error) { + collector := Collector{ + db: db, + } + + err := collector.getByJobID(ctx, jobID) + if err != nil { + return nil, err + } + + return &collector, nil +} + +func (c *Collector) getByJobID(ctx context.Context, jobID uuid.UUID) error { + dbJobID, err := database.ToDBUUID(jobID) + if err != nil { + return err + } + + dbCollector, err := c.db.GetCollectorFromJobID(ctx, dbJobID) + if err != nil { + return err + } + + c.ID = database.MustToUUID(dbCollector.ID) + c.MinCleanVersion = dbCollector.Mincleanversion + c.MinTextVersion = dbCollector.Mintextversion + + return nil +} diff --git a/internal/database/repository/models.go b/internal/database/repository/models.go index f551049d..9d93171c 100644 --- a/internal/database/repository/models.go +++ b/internal/database/repository/models.go @@ -110,7 +110,7 @@ type Result struct { ID pgtype.UUID Queryid pgtype.UUID Documentid pgtype.UUID - Value pgtype.Text + Value string Cleanversion int32 Textversion int32 Queryversion int32 diff --git a/internal/database/repository/result.sql.go b/internal/database/repository/result.sql.go index b1df57c3..fcacc232 100644 --- a/internal/database/repository/result.sql.go +++ b/internal/database/repository/result.sql.go @@ -11,8 +11,38 @@ import ( "github.com/jackc/pgx/v5/pgtype" ) +const listResultValuesByID = `-- name: ListResultValuesByID :many +SELECT id, queryId, value FROM results where id = ANY($1) +` + +type ListResultValuesByIDRow struct { + ID pgtype.UUID + Queryid pgtype.UUID + Value string +} + +func (q *Queries) ListResultValuesByID(ctx context.Context, id []pgtype.UUID) ([]ListResultValuesByIDRow, error) { + rows, err := q.db.Query(ctx, listResultValuesByID, id) + if err != nil { + return nil, err + } + defer rows.Close() + var items []ListResultValuesByIDRow + for rows.Next() { + var i ListResultValuesByIDRow + if err := rows.Scan(&i.ID, &i.Queryid, &i.Value); err != nil { + return nil, err + } + items = append(items, i) + } + if err := rows.Err(); err != nil { + return nil, err + } + return items, nil +} + const listResultsByDocumentID = `-- name: ListResultsByDocumentID :many -SELECT id, queryId, cleanVersion, textVersion, queryVersion FROM results where documentId = $1 and cleanVersion >= $2 and textVersion >= $3 +SELECT id, queryId, queryVersion FROM results where documentId = $1 and cleanVersion >= $2 and textVersion >= $3 ` type ListResultsByDocumentIDParams struct { @@ -24,8 +54,6 @@ type ListResultsByDocumentIDParams struct { type ListResultsByDocumentIDRow struct { ID pgtype.UUID Queryid pgtype.UUID - Cleanversion int32 - Textversion int32 Queryversion int32 } @@ -38,13 +66,7 @@ func (q *Queries) ListResultsByDocumentID(ctx context.Context, arg ListResultsBy var items []ListResultsByDocumentIDRow for rows.Next() { var i ListResultsByDocumentIDRow - if err := rows.Scan( - &i.ID, - &i.Queryid, - &i.Cleanversion, - &i.Textversion, - &i.Queryversion, - ); err != nil { + if err := rows.Scan(&i.ID, &i.Queryid, &i.Queryversion); err != nil { return nil, err } items = append(items, i) @@ -54,3 +76,30 @@ func (q *Queries) ListResultsByDocumentID(ctx context.Context, arg ListResultsBy } return items, nil } + +const setResult = `-- name: SetResult :exec +INSERT INTO results (id, queryId, documentId, value, cleanVersion, textVersion, queryVersion) VALUES ($1, $2, $3, $4, $5, $6, $7) +` + +type SetResultParams struct { + ID pgtype.UUID + Queryid pgtype.UUID + Documentid pgtype.UUID + Value string + Cleanversion int32 + Textversion int32 + Queryversion int32 +} + +func (q *Queries) SetResult(ctx context.Context, arg SetResultParams) error { + _, err := q.db.Exec(ctx, setResult, + arg.ID, + arg.Queryid, + arg.Documentid, + arg.Value, + arg.Cleanversion, + arg.Textversion, + arg.Queryversion, + ) + return err +} diff --git a/internal/database/uuid.go b/internal/database/uuid.go index 38ac53fd..6af0aaa8 100644 --- a/internal/database/uuid.go +++ b/internal/database/uuid.go @@ -1,6 +1,8 @@ package database import ( + "fmt" + "github.com/google/uuid" "github.com/jackc/pgx/v5/pgtype" ) @@ -14,3 +16,16 @@ func ToDBUUID(id uuid.UUID) (pgtype.UUID, error) { return dbID, nil } + +func MustToDBUUID(id uuid.UUID) pgtype.UUID { + dbID, err := ToDBUUID(id) + if err != nil { + panic(fmt.Sprint("id is not valid: ", id)) + } + + return dbID +} + +func MustToUUID(id pgtype.UUID) uuid.UUID { + return uuid.Must(uuid.FromBytes(id.Bytes[:])) +} diff --git a/internal/document/service.go b/internal/document/service.go index b2b6d39b..c124a0fc 100644 --- a/internal/document/service.go +++ b/internal/document/service.go @@ -2,49 +2,49 @@ package document import ( "context" + "gotemplate/internal/collector" "gotemplate/internal/database" "gotemplate/internal/database/repository" - "log" + "gotemplate/internal/query" "github.com/google/uuid" ) +type Document struct { + ID uuid.UUID `json:"id"` + JobID uuid.UUID `json:"jobId"` + Name string `json:"name"` + CleanVersion int32 `json:"cleanVersion"` + TextVersion int32 `json:"textVersion"` +} + type Service struct { db *repository.Queries } -type Document struct { - ID uuid.UUID `json:"id"` - JobID uuid.UUID `json:"jobId"` - Name string `json:"name"` -} - -type QueryType int - -const ( - QueryTypeJsonExtractor = iota -) - -type query struct { - ID uuid.UUID - Type QueryType - RequiredQueryID uuid.UUID - Version int -} - func New(ctx context.Context, db *repository.Queries) *Service { return &Service{ db: db, } } -func (s *Service) Sync(ctx context.Context, doc Document) error { - queries, err := s.getUnsyncedOrderedQueries(ctx, doc) +func (s *Service) Sync(ctx context.Context, doc *Document) error { + collector, err := collector.NewByJobId(ctx, s.db, doc.JobID) if err != nil { return err } - err = s.executeQueries(queries) + results, err := s.GetResults(ctx, doc.ID, collector) + if err != nil { + return err + } + + queue, err := query.NewQueue(ctx, s.db, collector, results, doc.ID, doc.CleanVersion, doc.TextVersion) + if err != nil { + return err + } + + err = queue.Execute(ctx) if err != nil { return err } @@ -52,65 +52,33 @@ func (s *Service) Sync(ctx context.Context, doc Document) error { return nil } -func (s *Service) getUnsyncedOrderedQueries(ctx context.Context, doc Document) ([]query, error) { - jobID, err := database.ToDBUUID(doc.JobID) - if err != nil { - return nil, err - } - - collector, err := s.db.GetCollectorFromJobID(ctx, jobID) - if err != nil { - return nil, err - } - - queries, err := s.db.GetCollectorQueries(ctx, collector.ID) - if err != nil { - return nil, err - } - - docID, err := database.ToDBUUID(doc.ID) +func (s *Service) GetResults(ctx context.Context, id uuid.UUID, coll *collector.Collector) (*[]query.Result, error) { + docID, err := database.ToDBUUID(id) if err != nil { return nil, err } results, err := s.db.ListResultsByDocumentID(ctx, repository.ListResultsByDocumentIDParams{ Documentid: docID, - Textversion: collector.Mintextversion, - Cleanversion: collector.Mincleanversion, + Textversion: coll.MinTextVersion, + Cleanversion: coll.MinCleanVersion, }) if err != nil { return nil, err } - queue := []query{} - - for _, query := range queries { - isSynced := false - for _, result := range results { - if result.Queryid != query.Queryid || int(result.Queryversion) != int(query.Queryversion.Int32) { - continue - } - isSynced = true - break - } - if isSynced { - continue - } - // TODO add query to beginning, add all dependents until one that is already in the queue + cleanResults := make([]query.Result, len(results)) + for index, dbResult := range results { + cleanResults[index] = *s.parseResult(&dbResult) } - return queue, nil + return &cleanResults, nil } -func (s *Service) executeQueries(queries []query) error { - // TODO - get all results? - - for _, query := range queries { - log.Print(query) - // TODO - get required field - // execute query - // store result fully +func (s *Service) parseResult(dbQuery *repository.ListResultsByDocumentIDRow) *query.Result { + return &query.Result{ + ID: database.MustToUUID(dbQuery.ID), + QueryID: database.MustToUUID(dbQuery.Queryid), + QueryVersion: dbQuery.Queryversion, } - - return nil } diff --git a/internal/query/database.go b/internal/query/database.go new file mode 100644 index 00000000..73262fef --- /dev/null +++ b/internal/query/database.go @@ -0,0 +1,38 @@ +package query + +import ( + "gotemplate/internal/database" + "gotemplate/internal/database/repository" +) + +func ParseDBQuery(dbQuery *repository.GetCollectorQueriesRow) (*Query, error) { + return &Query{ + ID: database.MustToUUID(dbQuery.Queryid), + Type: ParseDBQueryType(dbQuery.Type.Querytype), + RequiredQueryID: database.MustToUUID(dbQuery.Requiredqueryid), + Version: dbQuery.Queryversion.Int32, + }, nil +} + +func ParseDBQueryType(qType repository.Querytype) QueryType { + switch qType { + case repository.QuerytypeJsonExtractor: + return QueryTypeJsonExtractor + default: + return QueryTypeJsonExtractor + } +} + +func ToDBQueryType(t QueryType) repository.NullQuerytype { + var dbType repository.Querytype + + switch t { + case QueryTypeJsonExtractor: + dbType = repository.QuerytypeJsonExtractor + break + default: + dbType = repository.QuerytypeJsonExtractor + } + + return repository.NullQuerytype{Querytype: dbType, Valid: true} +} diff --git a/internal/query/service.go b/internal/query/service.go new file mode 100644 index 00000000..f055b990 --- /dev/null +++ b/internal/query/service.go @@ -0,0 +1,251 @@ +package query + +import ( + "context" + "gotemplate/internal/collector" + "gotemplate/internal/database" + "gotemplate/internal/database/repository" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgtype" +) + +type QueryType int + +const ( + QueryTypeJsonExtractor = iota +) + +type Query struct { + ID uuid.UUID + Type QueryType + RequiredQueryID uuid.UUID + Version int32 +} + +type Result struct { + ID uuid.UUID + QueryID uuid.UUID + QueryVersion int32 +} + +type ResultValue struct { + ID uuid.UUID + QueryID uuid.UUID + Value string +} + +type Queue struct { + unsyncedQueue *[]Query + collectorQueries *[]Query + results *[]Result + collector *collector.Collector + db *repository.Queries + cleanVersion int32 + textVersion int32 + documentId uuid.UUID +} + +func NewQueue(ctx context.Context, db *repository.Queries, coll *collector.Collector, results *[]Result, docId uuid.UUID, cleanVersion int32, textVersion int32) (*Queue, error) { + queue := Queue{ + db: db, + results: results, + collector: coll, + documentId: docId, + cleanVersion: cleanVersion, + textVersion: textVersion, + } + + err := queue.getCollectorQueries(ctx) + if err != nil { + return nil, err + } + + queue.getUnsyncedQueries() + + return &queue, nil +} + +func (q *Queue) GetQueue() []Query { + return *q.unsyncedQueue +} + +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) + } +} + +func (q *Queue) Execute(ctx context.Context) error { + if q.unsyncedQueue == nil { + return nil + } + + for _, query := range *q.unsyncedQueue { + err := q.executeQuery(ctx, query) + if err != nil { + return err + } + } + + q.unsyncedQueue = nil + + return nil +} + +func (q *Queue) executeQuery(ctx context.Context, query Query) error { + requiredQueryIDs := []uuid.UUID{} + for _, entry := range *q.collectorQueries { + if entry.ID == query.ID && entry.RequiredQueryID != uuid.Nil { + requiredQueryIDs = append(requiredQueryIDs, entry.RequiredQueryID) + } + } + + resultIDs := make([]pgtype.UUID, len(requiredQueryIDs)) + for index, id := range requiredQueryIDs { + var queryVersion int32 + for _, entry := range *q.collectorQueries { + if entry.ID == id { + queryVersion = entry.Version + break + } + } + + for _, entry := range *q.results { + if entry.QueryID == id && entry.QueryVersion == queryVersion { + resultIDs[index] = database.MustToDBUUID(entry.ID) + break + } + } + } + + values, err := q.db.ListResultValuesByID(ctx, resultIDs) + if err != nil { + return err + } + + cleanValues := make([]ResultValue, len(values)) + for _, result := range values { + cleanValue := q.parseResultValue(result) + cleanValues = append(cleanValues, cleanValue) + } + + err = q.getResult(ctx, query, &cleanValues) + if err != nil { + return err + } + + return nil +} + +func (q *Queue) parseResultValue(result repository.ListResultValuesByIDRow) ResultValue { + return ResultValue{ + ID: database.MustToUUID(result.ID), + QueryID: database.MustToUUID(result.Queryid), + Value: result.Value, + } +} + +func (q *Queue) getResult(ctx context.Context, query Query, resultValues *[]ResultValue) error { + // TODO - execute query + value := "the value" + + id := uuid.New() + + err := q.db.SetResult(ctx, repository.SetResultParams{ + ID: database.MustToDBUUID(id), + Queryid: database.MustToDBUUID(query.ID), + Documentid: database.MustToDBUUID(q.documentId), + Value: value, + Cleanversion: q.cleanVersion, + Textversion: q.textVersion, + Queryversion: query.Version, + }) + if err != nil { + return err + } + + *q.results = append(*q.results, Result{ + ID: id, + QueryID: query.ID, + QueryVersion: query.Version, + }) + + return nil +} diff --git a/test/unit/internal/document/sync_test.go b/test/unit/internal/document/sync_test.go index b0a9e078..33213e43 100644 --- a/test/unit/internal/document/sync_test.go +++ b/test/unit/internal/document/sync_test.go @@ -5,6 +5,7 @@ import ( "gotemplate/internal/database" "gotemplate/internal/database/repository" "gotemplate/internal/document" + "gotemplate/internal/sync" "testing" "github.com/google/uuid" @@ -23,30 +24,42 @@ func TestSync(t *testing.T) { defer db.Close(ctx) queries := repository.New(db) - svc := document.New(ctx, queries) + docSvc := document.New(ctx, queries) + svc := sync.New(ctx, queries, docSvc) + doc := document.Document{ ID: uuid.New(), JobID: uuid.New(), Name: "document_name", } + dbCollectorId, err := database.ToDBUUID(uuid.New()) + assert.Nil(t, err) dbJobID, err := database.ToDBUUID(doc.JobID) assert.Nil(t, err) - db.ExpectQuery("name: GetQueriesFromJobID :many").WithArgs(dbJobID). + minCleanVersion := int32(1) + minTextVersion := int32(1) + db.ExpectQuery("name: GetCollectorFromJobID :one").WithArgs(dbJobID). WillReturnRows( - pgxmock.NewRows([]string{"id", "type", "requiredQueryId", "queryVersion", "minCleanVersion", "minTextVersion"}). - AddRow(pgtype.UUID{}, repository.QuerytypeJsonExtractor, pgtype.UUID{}, int32(1), int32(1), int32(1)), + pgxmock.NewRows([]string{"id", "jobId", "minCleanVersion", "minTextVersion"}). + AddRow(dbCollectorId, dbJobID, minCleanVersion, minTextVersion), ) + assert.Nil(t, err) dbDocID, err := database.ToDBUUID(doc.ID) assert.Nil(t, err) - db.ExpectQuery("name: ListResultsByDocumentID :many").WithArgs(dbDocID). + db.ExpectQuery("name: ListResultsByDocumentID :many").WithArgs(dbDocID, minCleanVersion, minTextVersion). WillReturnRows( - pgxmock.NewRows([]string{"id", "queryId", "cleanVersion", "textVersion", "queryVersion"}). - AddRow(pgtype.UUID{}, pgtype.UUID{}, int32(1), int32(1), int32(1)), + pgxmock.NewRows([]string{"id", "queryId", "queryVersion"}). + AddRow(pgtype.UUID{}, pgtype.UUID{}, int32(1)), + ) + db.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorId). + WillReturnRows( + pgxmock.NewRows([]string{"collectorId", "queryId", "type", "requiredQueryId", "queryVersion"}). + AddRow(dbCollectorId, pgtype.UUID{}, repository.NullQuerytype{Querytype: repository.QuerytypeJsonExtractor, Valid: true}, pgtype.UUID{}, pgtype.Int4{Int32: int32(1), Valid: true}), ) - err = svc.Sync(ctx, doc) + err = svc.Sync(ctx, &doc) assert.Nil(t, err) } diff --git a/test/unit/internal/query/queue_test.go b/test/unit/internal/query/queue_test.go new file mode 100644 index 00000000..f8c38c82 --- /dev/null +++ b/test/unit/internal/query/queue_test.go @@ -0,0 +1,125 @@ +package document_test + +import ( + "context" + "gotemplate/internal/collector" + "gotemplate/internal/database" + "gotemplate/internal/database/repository" + "gotemplate/internal/query" + "testing" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgtype" + "github.com/pashagolub/pgxmock/v3" + "github.com/stretchr/testify/assert" +) + +func TestQueue(t *testing.T) { + ctx := context.Background() + + db, err := pgxmock.NewConn() + if err != nil { + t.Fatalf("failed to open pgxmock database: %v", err) + } + defer db.Close(ctx) + + queries := repository.New(db) + jobID := uuid.New() + dbJobID := database.MustToDBUUID(jobID) + collectorID := uuid.New() + dbCollectorID := database.MustToDBUUID(collectorID) + + db.ExpectQuery("name: GetCollectorFromJobID :one").WithArgs(dbJobID). + WillReturnRows( + pgxmock.NewRows([]string{"id", "jobId", "minCleanVersion", "minTextVersion"}). + AddRow(dbCollectorID, dbJobID, int32(1), int32(1)), + ) + assert.Nil(t, err) + + coll, err := collector.NewByJobId(ctx, queries, jobID) + assert.Nil(t, err) + + queryOneID := uuid.New() + queryOneVersion := int32(1) + queryTwoID := uuid.New() + queryTwoVersion := int32(2) + queryThreeID := uuid.New() + queryThreeVersion := int32(3) + queryFourID := uuid.New() + queryFourVersion := int32(4) + queryFiveID := uuid.New() + queryFiveVersion := int32(5) + querySixID := uuid.New() + querySixVersion := int32(6) + collectorQueries := []query.Query{ + {ID: queryOneID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: uuid.Nil, Version: queryOneVersion}, + {ID: queryTwoID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: queryOneID, Version: queryTwoVersion}, + {ID: queryThreeID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: queryOneID, Version: queryThreeVersion}, + {ID: queryThreeID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: queryTwoID, Version: queryThreeVersion}, + {ID: queryFourID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: uuid.Nil, Version: queryFourVersion}, + {ID: queryFiveID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: querySixID, Version: queryFiveVersion}, + {ID: querySixID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: uuid.Nil, Version: querySixVersion}, + } + + rows := pgxmock.NewRows([]string{"collectorId", "queryId", "type", "requiredQueryId", "queryVersion"}) + for _, q := range collectorQueries { + dbID := database.MustToDBUUID(q.ID) + dbReqID := database.MustToDBUUID(q.RequiredQueryID) + rows = rows. + AddRow(dbCollectorID, dbID, query.ToDBQueryType(q.Type), dbReqID, pgtype.Int4{Int32: int32(q.Version), Valid: true}) + } + + db.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID).WillReturnRows(rows) + assert.Nil(t, err) + + results := []query.Result{ + {ID: uuid.New(), QueryID: queryFourID, QueryVersion: queryFourVersion}, + {ID: uuid.New(), QueryID: querySixID, QueryVersion: querySixVersion - 1}, + {ID: uuid.New(), QueryID: queryOneID, QueryVersion: queryOneVersion - 1}, + } + + expectedQueries := []query.Query{ + {ID: querySixID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: uuid.Nil, Version: querySixVersion}, + {ID: queryFiveID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: querySixID, Version: queryFiveVersion}, + {ID: queryThreeID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: queryTwoID, Version: queryThreeVersion}, + {ID: queryTwoID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: queryOneID, Version: queryTwoVersion}, + {ID: queryOneID, Type: query.QueryTypeJsonExtractor, RequiredQueryID: uuid.Nil, Version: queryOneVersion}, + } + + docID := uuid.New() + cleanVersion := int32(1) + textVersion := int32(1) + + q, err := query.NewQueue(ctx, queries, coll, &results, docID, cleanVersion, textVersion) + assert.Nil(t, err) + assert.Equal(t, expectedQueries, q.GetQueue()) + + db.ExpectQuery("name: ListResultValuesByID :many").WithArgs([]pgtype.UUID{}).WillReturnRows(rows) + assert.Nil(t, err) + db.ExpectExec("name: SetResult :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID(querySixID), database.MustToDBUUID(docID), pgxmock.AnyArg(), cleanVersion, textVersion, querySixVersion). + WillReturnResult(pgxmock.NewResult("", 1)) + assert.Nil(t, err) + db.ExpectQuery("name: ListResultValuesByID :many").WithArgs(pgxmock.AnyArg()).WillReturnRows(rows) + assert.Nil(t, err) + db.ExpectExec("name: SetResult :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID(queryFiveID), database.MustToDBUUID(docID), pgxmock.AnyArg(), cleanVersion, textVersion, queryFiveVersion). + WillReturnResult(pgxmock.NewResult("", 1)) + assert.Nil(t, err) + db.ExpectQuery("name: ListResultValuesByID :many").WithArgs(pgxmock.AnyArg()).WillReturnRows(rows) + assert.Nil(t, err) + db.ExpectExec("name: SetResult :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID(queryThreeID), database.MustToDBUUID(docID), pgxmock.AnyArg(), cleanVersion, textVersion, queryThreeVersion). + WillReturnResult(pgxmock.NewResult("", 1)) + assert.Nil(t, err) + db.ExpectQuery("name: ListResultValuesByID :many").WithArgs(pgxmock.AnyArg()).WillReturnRows(rows) + assert.Nil(t, err) + db.ExpectExec("name: SetResult :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID(queryTwoID), database.MustToDBUUID(docID), pgxmock.AnyArg(), cleanVersion, textVersion, queryTwoVersion). + WillReturnResult(pgxmock.NewResult("", 1)) + assert.Nil(t, err) + db.ExpectQuery("name: ListResultValuesByID :many").WithArgs(pgxmock.AnyArg()).WillReturnRows(rows) + assert.Nil(t, err) + db.ExpectExec("name: SetResult :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID(queryOneID), database.MustToDBUUID(docID), pgxmock.AnyArg(), cleanVersion, textVersion, queryOneVersion). + WillReturnResult(pgxmock.NewResult("", 1)) + assert.Nil(t, err) + + err = q.Execute(ctx) + assert.Nil(t, err) +}