Merged in feature/shorttestsanddirtidy (pull request #26)

Add short tests and Tidy internal directories

* complete the tasks
This commit is contained in:
Michael McGuinness
2025-01-17 12:00:32 +00:00
parent b453f6cb23
commit fa95d733ca
84 changed files with 171 additions and 101 deletions
+94
View File
@@ -0,0 +1,94 @@
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.GetCollectorQueries(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
}
for _, id := range qu.RequiredQueryIDs {
if entry.ID == id {
requiredIndex = index
break
}
}
}
}
for _, entry := range q.collectorQueries {
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)
}
}
+139
View File
@@ -0,0 +1,139 @@
package queryqueue
import (
"context"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/job/collector"
queryprocessor "queryorchestration/internal/query/processor"
"queryorchestration/internal/query/result"
"testing"
"github.com/google/uuid"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
)
func TestGetUnsyncedQueries(t *testing.T) {
queryOne := &queryprocessor.Query{
ID: uuid.New(),
Version: int32(1),
}
svc := Queue{
collectorQueries: []*queryprocessor.Query{
queryOne,
},
results: []*result.Result{
{ID: uuid.New(), QueryID: queryOne.ID, QueryVersion: queryOne.Version},
},
}
svc.getUnsyncedQueries()
assert.EqualExportedValues(t, []*queryprocessor.Query(nil), svc.unsyncedQueue)
svc.results = []*result.Result{}
svc.unsyncedQueue = []*queryprocessor.Query{}
svc.getUnsyncedQueries()
assert.EqualExportedValues(t, []*queryprocessor.Query{queryOne}, svc.unsyncedQueue)
}
func TestGetCollectorQueries(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
queries := repository.New(pool)
db := &database.Connection{
Queries: queries,
Pool: pool,
}
svc := Queue{
db: db,
collector: &collector.Collector{
ID: uuid.New(),
},
}
dbCollectorID := database.MustToDBUUID(svc.collector.ID)
collectorQueries := []*queryprocessor.Query{
{ID: uuid.New(), Type: queryprocessor.TypeContextFull, RequiredQueryIDs: []uuid.UUID{}, Version: int32(1)},
}
rows := pgxmock.NewRows([]string{"collectorId", "queryId", "type", "queryVersion", "requiredIds"})
for _, q := range collectorQueries {
dbID := database.MustToDBUUID(q.ID)
dbReqIDs := database.MustToDBUUIDArray(q.RequiredQueryIDs)
ty, err := queryprocessor.ToDBNullQueryType(q.Type)
assert.Nil(t, err)
rows = rows.
AddRow(dbCollectorID, dbID, ty, &q.Version, dbReqIDs)
}
pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID).WillReturnRows(rows)
err = svc.getCollectorQueries(ctx)
assert.Nil(t, err)
assert.EqualExportedValues(t, collectorQueries, svc.collectorQueries)
}
func TestIsQuerySynced(t *testing.T) {
query := &queryprocessor.Query{
ID: uuid.New(),
Version: int32(1),
}
svc := Queue{
results: []*result.Result{
{QueryID: query.ID, QueryVersion: query.Version},
},
}
isSynced := svc.isQuerySynced(query)
assert.True(t, isSynced)
}
func TestIsQuerySyncedNoResult(t *testing.T) {
query := &queryprocessor.Query{
ID: uuid.New(),
Version: int32(1),
}
svc := Queue{
results: []*result.Result{},
}
isSynced := svc.isQuerySynced(query)
assert.False(t, isSynced)
}
func TestIsQuerySyncedOldResult(t *testing.T) {
query := &queryprocessor.Query{
ID: uuid.New(),
Version: int32(1),
}
svc := Queue{
results: []*result.Result{
{QueryID: query.ID, QueryVersion: query.Version - 1},
},
}
isSynced := svc.isQuerySynced(query)
assert.False(t, isSynced)
}
func TestIsQuerySyncedNewResult(t *testing.T) {
query := &queryprocessor.Query{
ID: uuid.New(),
Version: int32(1),
}
svc := Queue{
results: []*result.Result{
{QueryID: query.ID, QueryVersion: query.Version + 1},
},
}
isSynced := svc.isQuerySynced(query)
assert.False(t, isSynced)
}
+69
View File
@@ -0,0 +1,69 @@
package queryqueue
import (
"context"
"queryorchestration/internal/database"
queryprocessor "queryorchestration/internal/query/processor"
"queryorchestration/internal/query/result"
"github.com/jackc/pgx/v5/pgtype"
)
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, qu *queryprocessor.Query) error {
resultIDs := make([]pgtype.UUID, len(qu.RequiredQueryIDs))
for index, id := range qu.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.Queries.ListResultValuesByID(ctx, resultIDs)
if err != nil {
return err
}
cleanValues := make([]result.Value, len(values))
for index, r := range values {
cleanValue, err := q.getResultValue(r)
if err != nil {
return err
}
cleanValues[index] = cleanValue
}
err = q.setResult(ctx, qu, cleanValues)
if err != nil {
return err
}
return nil
}
+179
View File
@@ -0,0 +1,179 @@
package queryqueue_test
import (
"context"
"fmt"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/job/collector"
queryprocessor "queryorchestration/internal/query/processor"
queryqueue "queryorchestration/internal/query/queue"
"queryorchestration/internal/query/result"
"testing"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgtype"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
)
func TestExecute(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
queries := repository.New(pool)
db := &database.Connection{
Queries: queries,
Pool: pool,
}
jobID := uuid.New()
dbJobID := database.MustToDBUUID(jobID)
collectorID := uuid.New()
dbCollectorID := database.MustToDBUUID(collectorID)
pool.ExpectQuery("name: GetCollectorFromJobID :one").WithArgs(dbJobID).
WillReturnRows(
pgxmock.NewRows([]string{"id", "jobId", "minCleanVersion", "minTextVersion"}).
AddRow(dbCollectorID, dbJobID, int32(1), int32(1)),
)
coll, err := collector.NewByJobId(ctx, db, 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)
contextID := uuid.New()
contextVersion := int32(1)
collectorQueries := []queryprocessor.Query{
{ID: contextID, Type: queryprocessor.TypeContextFull, RequiredQueryIDs: []uuid.UUID{}, Version: contextVersion},
{ID: queryOneID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: queryOneVersion},
{ID: queryTwoID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{queryOneID}, Version: queryTwoVersion},
{ID: queryThreeID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{queryOneID}, Version: queryThreeVersion},
{ID: queryFourID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: queryFourVersion},
{ID: queryFiveID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{querySixID}, Version: queryFiveVersion},
{ID: querySixID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: querySixVersion},
}
rows := pgxmock.NewRows([]string{"collectorId", "queryId", "type", "queryVersion", "requiredIds"})
for _, q := range collectorQueries {
dbID := database.MustToDBUUID(q.ID)
dbReqIDs := database.MustToDBUUIDArray(q.RequiredQueryIDs)
ty, err := queryprocessor.ToDBNullQueryType(q.Type)
assert.Nil(t, err)
rows = rows.
AddRow(dbCollectorID, dbID, ty, &q.Version, dbReqIDs)
}
pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID).WillReturnRows(rows)
contextResultID := uuid.New()
results := []*result.Result{
{ID: contextResultID, QueryID: contextID, QueryVersion: contextVersion},
{ID: uuid.New(), QueryID: queryFourID, QueryVersion: queryFourVersion},
{ID: uuid.New(), QueryID: querySixID, QueryVersion: querySixVersion - 1},
{ID: uuid.New(), QueryID: queryOneID, QueryVersion: queryOneVersion - 1},
}
expectedQueries := []*queryprocessor.Query{
{ID: querySixID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: querySixVersion},
{ID: queryFiveID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{querySixID}, Version: queryFiveVersion},
{ID: queryOneID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: queryOneVersion},
{ID: queryThreeID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{queryOneID}, Version: queryThreeVersion},
{ID: queryTwoID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{queryOneID}, Version: queryTwoVersion},
}
docID := uuid.New()
cleanVersion := int32(1)
textVersion := int32(1)
q, err := queryqueue.New(ctx, db, coll, results, docID, cleanVersion, textVersion)
assert.Nil(t, err)
assert.Equal(t, expectedQueries, q.GetQueue())
keyLayerOne := "key"
keyLayerTwo := "key5"
valueLayerTwo := "value"
valueLayerOne := fmt.Sprintf("{\"%s\":\"%s\"}", keyLayerTwo, valueLayerTwo)
valueContext := fmt.Sprintf("{\"%s\":%s}", keyLayerOne, valueLayerOne)
pool.ExpectQuery("name: ListResultValuesByID :many").WithArgs([]pgtype.UUID{database.MustToDBUUID(contextResultID)}).WillReturnRows(
pgxmock.NewRows([]string{"id", "queryId", "value"}).
AddRow(database.MustToDBUUID(contextResultID), database.MustToDBUUID(contextID), valueContext),
)
pool.ExpectQuery("name: GetQueryConfig :one").WithArgs(database.MustToDBUUID(querySixID), querySixVersion).WillReturnRows(
pgxmock.NewRows([]string{"id", "config"}).
AddRow(pgtype.UUID{}, []byte(fmt.Sprintf("{\"path\":\"%s\"}", keyLayerOne))),
)
pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(querySixID), database.MustToDBUUID(docID), valueLayerOne, cleanVersion, textVersion, querySixVersion).
WillReturnRows(
pgxmock.NewRows([]string{"id"}).AddRow(pgtype.UUID{}),
)
pool.ExpectQuery("name: ListResultValuesByID :many").WithArgs(pgxmock.AnyArg()).WillReturnRows(
pgxmock.NewRows([]string{"id", "queryId", "value"}).
AddRow(pgtype.UUID{}, database.MustToDBUUID(querySixID), valueLayerOne),
)
pool.ExpectQuery("name: GetQueryConfig :one").WithArgs(database.MustToDBUUID(queryFiveID), queryFiveVersion).WillReturnRows(
pgxmock.NewRows([]string{"id", "config"}).
AddRow(pgtype.UUID{}, []byte(fmt.Sprintf("{\"path\":\"%s\"}", keyLayerTwo))),
)
pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(queryFiveID), database.MustToDBUUID(docID), valueLayerTwo, cleanVersion, textVersion, queryFiveVersion).
WillReturnRows(
pgxmock.NewRows([]string{"id"}).AddRow(pgtype.UUID{}),
)
pool.ExpectQuery("name: ListResultValuesByID :many").WithArgs(pgxmock.AnyArg()).WillReturnRows(
pgxmock.NewRows([]string{"id", "queryId", "value"}).
AddRow(database.MustToDBUUID(contextResultID), database.MustToDBUUID(contextID), valueContext),
)
pool.ExpectQuery("name: GetQueryConfig :one").WithArgs(database.MustToDBUUID(queryOneID), queryOneVersion).WillReturnRows(
pgxmock.NewRows([]string{"id", "config"}).
AddRow(pgtype.UUID{}, []byte(fmt.Sprintf("{\"path\":\"%s\"}", keyLayerOne))),
)
pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(queryOneID), database.MustToDBUUID(docID), valueLayerOne, cleanVersion, textVersion, queryOneVersion).
WillReturnRows(
pgxmock.NewRows([]string{"id"}).AddRow(pgtype.UUID{}),
)
pool.ExpectQuery("name: ListResultValuesByID :many").WithArgs(pgxmock.AnyArg()).WillReturnRows(
pgxmock.NewRows([]string{"id", "queryId", "value"}).
AddRow(pgtype.UUID{}, database.MustToDBUUID(queryOneID), valueLayerOne),
)
pool.ExpectQuery("name: GetQueryConfig :one").WithArgs(database.MustToDBUUID(queryThreeID), queryThreeVersion).WillReturnRows(
pgxmock.NewRows([]string{"id", "config"}).
AddRow(pgtype.UUID{}, []byte(fmt.Sprintf("{\"path\":\"%s\"}", keyLayerTwo))),
)
pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(queryThreeID), database.MustToDBUUID(docID), valueLayerTwo, cleanVersion, textVersion, queryThreeVersion).
WillReturnRows(
pgxmock.NewRows([]string{"id"}).AddRow(pgtype.UUID{}),
)
pool.ExpectQuery("name: ListResultValuesByID :many").WithArgs(pgxmock.AnyArg()).WillReturnRows(
pgxmock.NewRows([]string{"id", "queryId", "value"}).
AddRow(pgtype.UUID{}, database.MustToDBUUID(queryOneID), valueLayerOne),
)
pool.ExpectQuery("name: GetQueryConfig :one").WithArgs(database.MustToDBUUID(queryTwoID), queryTwoVersion).WillReturnRows(
pgxmock.NewRows([]string{"id", "config"}).
AddRow(pgtype.UUID{}, []byte(fmt.Sprintf("{\"path\":\"%s\"}", keyLayerTwo))),
)
pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(queryTwoID), database.MustToDBUUID(docID), valueLayerTwo, cleanVersion, textVersion, queryTwoVersion).
WillReturnRows(
pgxmock.NewRows([]string{"id"}).AddRow(pgtype.UUID{}),
)
err = q.Execute(ctx)
assert.Nil(t, err)
}
@@ -0,0 +1,94 @@
package queryqueue
import (
"context"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
queryprocessor "queryorchestration/internal/query/processor"
"testing"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgtype"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
)
func TestExecute(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
queries := repository.New(pool)
db := &database.Connection{
Queries: queries,
Pool: pool,
}
expectedQueries := []*queryprocessor.Query{
{ID: uuid.New(), Type: queryprocessor.TypeContextFull, RequiredQueryIDs: []uuid.UUID{}, Version: int32(1)},
{ID: uuid.New(), Type: queryprocessor.TypeContextFull, RequiredQueryIDs: []uuid.UUID{}, Version: int32(2)},
}
q := &Queue{
unsyncedQueue: expectedQueries,
db: db,
documentId: uuid.New(),
cleanVersion: int32(1),
textVersion: int32(2),
}
pool.ExpectQuery("name: ListResultValuesByID :many").WithArgs([]pgtype.UUID{}).WillReturnRows(
pgxmock.NewRows([]string{"id", "queryId", "value"}),
)
pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(expectedQueries[0].ID), database.MustToDBUUID(q.documentId), pgxmock.AnyArg(), q.cleanVersion, q.textVersion, expectedQueries[0].Version).
WillReturnRows(
pgxmock.NewRows([]string{"id"}).AddRow(pgtype.UUID{}),
)
pool.ExpectQuery("name: ListResultValuesByID :many").WithArgs([]pgtype.UUID{}).WillReturnRows(
pgxmock.NewRows([]string{"id", "queryId", "value"}),
)
pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(expectedQueries[1].ID), database.MustToDBUUID(q.documentId), pgxmock.AnyArg(), q.cleanVersion, q.textVersion, expectedQueries[1].Version).
WillReturnRows(
pgxmock.NewRows([]string{"id"}).AddRow(pgtype.UUID{}),
)
err = q.Execute(ctx)
assert.Nil(t, err)
}
func TestExecuteQuery(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
queries := repository.New(pool)
db := &database.Connection{
Queries: queries,
Pool: pool,
}
qu := &queryprocessor.Query{ID: uuid.New(), Type: queryprocessor.TypeContextFull, RequiredQueryIDs: []uuid.UUID{}, Version: int32(1)}
q := &Queue{
db: db,
documentId: uuid.New(),
cleanVersion: int32(1),
textVersion: int32(2),
}
pool.ExpectQuery("name: ListResultValuesByID :many").WithArgs([]pgtype.UUID{}).WillReturnRows(
pgxmock.NewRows([]string{"id", "queryId", "value"}),
)
pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(qu.ID), database.MustToDBUUID(q.documentId), pgxmock.AnyArg(), q.cleanVersion, q.textVersion, qu.Version).
WillReturnRows(
pgxmock.NewRows([]string{"id"}).AddRow(pgtype.UUID{}),
)
err = q.executeQuery(ctx, qu)
assert.Nil(t, err)
}
+73
View File
@@ -0,0 +1,73 @@
package queryqueue
import (
"context"
"fmt"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
queryprocessor "queryorchestration/internal/query/processor"
"queryorchestration/internal/query/result"
contextfull "queryorchestration/internal/query/types/contextFull"
jsonextractor "queryorchestration/internal/query/types/jsonExtractor"
)
func (q *Queue) setResult(ctx context.Context, qu *queryprocessor.Query, resultValues []result.Value) error {
processor, err := q.getProcessor(qu.Type)
if err != nil {
return err
}
value, err := processor.Process(ctx, qu, resultValues)
if err != nil {
return err
}
id, err := result.Store(ctx, q.db.Queries, &result.ResultStore{
QueryID: qu.ID,
DocumentID: q.documentId,
Value: value,
CleanVersion: q.cleanVersion,
TextVersion: q.textVersion,
QueryVersion: qu.Version,
})
if err != nil {
return err
}
q.results = append(q.results, &result.Result{
ID: id,
QueryID: qu.ID,
QueryVersion: qu.Version,
})
return nil
}
func (q *Queue) getProcessor(queryType queryprocessor.Type) (queryprocessor.Processor, error) {
switch queryType {
case queryprocessor.TypeJsonExtractor:
return jsonextractor.NewExtractor(q.db), nil
case queryprocessor.TypeContextFull:
return contextfull.NewExtractor(), nil
default:
return nil, fmt.Errorf("attempting to process invalid query type")
}
}
func (q *Queue) getResultValue(res *repository.ListResultValuesByIDRow) (result.Value, error) {
var queryType queryprocessor.Type
for _, qu := range q.collectorQueries {
if qu.ID == database.MustToUUID(res.Queryid) {
queryType = qu.Type
}
}
switch queryType {
case queryprocessor.TypeJsonExtractor:
return jsonextractor.NewResult(res.Value), nil
case queryprocessor.TypeContextFull:
return contextfull.NewResult(res.Value), nil
default:
return nil, fmt.Errorf("attempting to process invalid query type")
}
}
@@ -0,0 +1,74 @@
package queryqueue
import (
"context"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
queryprocessor "queryorchestration/internal/query/processor"
"queryorchestration/internal/query/result"
"testing"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgtype"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
)
func TestSetResult(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
queries := repository.New(pool)
db := &database.Connection{
Queries: queries,
Pool: pool,
}
q := &Queue{
db: db,
documentId: uuid.New(),
cleanVersion: int32(1),
textVersion: int32(2),
}
qu := &queryprocessor.Query{ID: uuid.New(), Type: queryprocessor.TypeContextFull, RequiredQueryIDs: []uuid.UUID{}, Version: int32(1)}
resultValues := []result.Value{}
pool.ExpectQuery("name: SetResult :one").WithArgs(database.MustToDBUUID(qu.ID), database.MustToDBUUID(q.documentId), pgxmock.AnyArg(), q.cleanVersion, q.textVersion, qu.Version).
WillReturnRows(
pgxmock.NewRows([]string{"id"}).AddRow(pgtype.UUID{}),
)
err = q.setResult(ctx, qu, resultValues)
assert.Nil(t, err)
}
func TestGetProcessor(t *testing.T) {
q := Queue{}
qType := queryprocessor.Type(queryprocessor.TypeContextFull)
processor, err := q.getProcessor(qType)
assert.Nil(t, err)
assert.NotNil(t, processor)
}
func TestGetResultValue(t *testing.T) {
result := &repository.ListResultValuesByIDRow{
Queryid: database.MustToDBUUID(uuid.New()),
Value: "EXAMPLE_VALUE",
}
q := &Queue{
collectorQueries: []*queryprocessor.Query{
{ID: database.MustToUUID(result.Queryid), Type: queryprocessor.TypeContextFull},
},
}
value, err := q.getResultValue(result)
assert.Nil(t, err)
assert.NotNil(t, value)
}
+46
View File
@@ -0,0 +1,46 @@
package queryqueue
import (
"context"
"queryorchestration/internal/database"
"queryorchestration/internal/job/collector"
queryprocessor "queryorchestration/internal/query/processor"
"queryorchestration/internal/query/result"
"github.com/google/uuid"
)
type Queue struct {
unsyncedQueue []*queryprocessor.Query
collectorQueries []*queryprocessor.Query
results []*result.Result
collector *collector.Collector
db *database.Connection
cleanVersion int32
textVersion int32
documentId uuid.UUID
}
func New(ctx context.Context, db *database.Connection, coll *collector.Collector, results []*result.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() []*queryprocessor.Query {
return q.unsyncedQueue
}
+144
View File
@@ -0,0 +1,144 @@
package queryqueue_test
import (
"context"
"errors"
"queryorchestration/internal/database"
"queryorchestration/internal/database/repository"
"queryorchestration/internal/job/collector"
queryprocessor "queryorchestration/internal/query/processor"
queryqueue "queryorchestration/internal/query/queue"
"queryorchestration/internal/query/result"
"testing"
"github.com/google/uuid"
"github.com/pashagolub/pgxmock/v3"
"github.com/stretchr/testify/assert"
)
func TestService(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
queries := repository.New(pool)
db := &database.Connection{
Queries: queries,
Pool: pool,
}
jobID := uuid.New()
dbJobID := database.MustToDBUUID(jobID)
collectorID := uuid.New()
dbCollectorID := database.MustToDBUUID(collectorID)
pool.ExpectQuery("name: GetCollectorFromJobID :one").WithArgs(dbJobID).
WillReturnRows(
pgxmock.NewRows([]string{"id", "jobId", "minCleanVersion", "minTextVersion"}).
AddRow(dbCollectorID, dbJobID, int32(1), int32(1)),
)
coll, err := collector.NewByJobId(ctx, db, 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)
contextID := uuid.New()
contextVersion := int32(1)
collectorQueries := []queryprocessor.Query{
{ID: contextID, Type: queryprocessor.TypeContextFull, RequiredQueryIDs: []uuid.UUID{}, Version: contextVersion},
{ID: queryOneID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: queryOneVersion},
{ID: queryTwoID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{queryOneID}, Version: queryTwoVersion},
{ID: queryThreeID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{queryOneID}, Version: queryThreeVersion},
{ID: queryFourID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: queryFourVersion},
{ID: queryFiveID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{querySixID}, Version: queryFiveVersion},
{ID: querySixID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: querySixVersion},
}
rows := pgxmock.NewRows([]string{"collectorId", "queryId", "type", "queryVersion", "requiredIds"})
for _, q := range collectorQueries {
dbID := database.MustToDBUUID(q.ID)
dbReqIDs := database.MustToDBUUIDArray(q.RequiredQueryIDs)
ty, err := queryprocessor.ToDBNullQueryType(q.Type)
assert.Nil(t, err)
rows = rows.
AddRow(dbCollectorID, dbID, ty, &q.Version, dbReqIDs)
}
pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID).WillReturnRows(rows)
contextResultID := uuid.New()
results := []*result.Result{
{ID: contextResultID, QueryID: contextID, QueryVersion: contextVersion},
{ID: uuid.New(), QueryID: queryFourID, QueryVersion: queryFourVersion},
{ID: uuid.New(), QueryID: querySixID, QueryVersion: querySixVersion - 1},
{ID: uuid.New(), QueryID: queryOneID, QueryVersion: queryOneVersion - 1},
}
expectedQueries := []*queryprocessor.Query{
{ID: querySixID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: querySixVersion},
{ID: queryFiveID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{querySixID}, Version: queryFiveVersion},
{ID: queryOneID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{contextID}, Version: queryOneVersion},
{ID: queryThreeID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{queryOneID}, Version: queryThreeVersion},
{ID: queryTwoID, Type: queryprocessor.TypeJsonExtractor, RequiredQueryIDs: []uuid.UUID{queryOneID}, Version: queryTwoVersion},
}
docID := uuid.New()
cleanVersion := int32(1)
textVersion := int32(1)
q, err := queryqueue.New(ctx, db, coll, results, docID, cleanVersion, textVersion)
assert.Nil(t, err)
assert.Equal(t, expectedQueries, q.GetQueue())
}
func TestQueueFail(t *testing.T) {
ctx := context.Background()
pool, err := pgxmock.NewPool()
if err != nil {
t.Fatalf("failed to open pgxmock database: %v", err)
}
queries := repository.New(pool)
db := &database.Connection{
Queries: queries,
Pool: pool,
}
jobID := uuid.New()
dbJobID := database.MustToDBUUID(jobID)
collectorID := uuid.New()
dbCollectorID := database.MustToDBUUID(collectorID)
pool.ExpectQuery("name: GetCollectorFromJobID :one").WithArgs(dbJobID).
WillReturnRows(
pgxmock.NewRows([]string{"id", "jobId", "minCleanVersion", "minTextVersion"}).
AddRow(dbCollectorID, dbJobID, int32(1), int32(1)),
)
coll, err := collector.NewByJobId(ctx, db, jobID)
assert.Nil(t, err)
dbErr := "database failure"
pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID).
WillReturnError(errors.New(dbErr))
results := []*result.Result{}
docID := uuid.New()
cleanVersion := int32(1)
textVersion := int32(1)
_, err = queryqueue.New(ctx, db, coll, results, docID, cleanVersion, textVersion)
assert.EqualError(t, err, dbErr)
}