somecleanup
This commit is contained in:
@@ -2,7 +2,6 @@ package controllers
|
||||
|
||||
import (
|
||||
"context"
|
||||
"queryorchestration/api/grpc/spec"
|
||||
serviceinterfaces "queryorchestration/api/serviceInterfaces"
|
||||
"queryorchestration/internal/query"
|
||||
queryprocessor "queryorchestration/internal/queryProcessor"
|
||||
@@ -41,9 +40,9 @@ func (s *QueryController) List(ctx context.Context, req *serviceinterfaces.Query
|
||||
return nil, err
|
||||
}
|
||||
|
||||
outQueries := make([]*serviceinterfaces.Query, len(*queries))
|
||||
for index, query := range *queries {
|
||||
outQueries[index] = ParseQuery(&query)
|
||||
outQueries := make([]*serviceinterfaces.Query, len(queries))
|
||||
for index, query := range queries {
|
||||
outQueries[index] = ParseQuery(query)
|
||||
}
|
||||
|
||||
return &serviceinterfaces.Queries{
|
||||
@@ -105,7 +104,7 @@ func (s *QueryController) Update(ctx context.Context, req *serviceinterfaces.Que
|
||||
return &emptypb.Empty{}, nil
|
||||
}
|
||||
|
||||
func (s *QueryController) Deprecate(ctx context.Context, req *spec.IdMessage) (*emptypb.Empty, error) {
|
||||
func (s *QueryController) Deprecate(ctx context.Context, req *serviceinterfaces.IdMessage) (*emptypb.Empty, error) {
|
||||
id, err := uuid.Parse(req.GetId())
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
||||
@@ -8,11 +8,11 @@ type Result struct {
|
||||
value string
|
||||
}
|
||||
|
||||
func NewResult(value string) *Result {
|
||||
return &Result{value}
|
||||
func NewResult(value string) Result {
|
||||
return Result{value}
|
||||
}
|
||||
|
||||
func (r *Result) GetValue(ctx context.Context) (string, error) {
|
||||
func (r Result) GetValue(ctx context.Context) (string, error) {
|
||||
// TODO - get value from s3
|
||||
return r.value, nil
|
||||
}
|
||||
|
||||
@@ -14,8 +14,8 @@ func NewExtractor() Extractor {
|
||||
return Extractor{}
|
||||
}
|
||||
|
||||
func (e Extractor) Process(ctx context.Context, query queryprocessor.Query, values *[]result.Value) (string, error) {
|
||||
if len(*values) > 0 {
|
||||
func (e Extractor) Process(ctx context.Context, query *queryprocessor.Query, values []result.Value) (string, error) {
|
||||
if len(values) > 0 {
|
||||
return "", fmt.Errorf("no requirements expected")
|
||||
}
|
||||
// TODO
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
)
|
||||
|
||||
func MustToDBUUIDArray(ids []uuid.UUID) []pgtype.UUID {
|
||||
dbIDs := make([]pgtype.UUID, len(ids))
|
||||
for index, id := range ids {
|
||||
dbIDs[index] = MustToDBUUID(id)
|
||||
}
|
||||
|
||||
return dbIDs
|
||||
}
|
||||
|
||||
func MustToDBUUID(id uuid.UUID) pgtype.UUID {
|
||||
var dbID pgtype.UUID
|
||||
dbID.Scan(id.String())
|
||||
return dbID
|
||||
}
|
||||
|
||||
func MustToUUID(id pgtype.UUID) uuid.UUID {
|
||||
return uuid.Must(uuid.FromBytes(id.Bytes[:]))
|
||||
}
|
||||
|
||||
func MustToUUIDArray(dbIDs []pgtype.UUID) []uuid.UUID {
|
||||
ids := make([]uuid.UUID, len(dbIDs))
|
||||
for index, id := range dbIDs {
|
||||
ids[index] = MustToUUID(id)
|
||||
}
|
||||
|
||||
return ids
|
||||
}
|
||||
@@ -1,16 +0,0 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
)
|
||||
|
||||
func MustToDBUUID(id uuid.UUID) pgtype.UUID {
|
||||
var dbID pgtype.UUID
|
||||
dbID.Scan(id.String())
|
||||
return dbID
|
||||
}
|
||||
|
||||
func MustToUUID(id pgtype.UUID) uuid.UUID {
|
||||
return uuid.Must(uuid.FromBytes(id.Bytes[:]))
|
||||
}
|
||||
@@ -53,7 +53,7 @@ func (s *Service) Sync(ctx context.Context, doc *Document) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Service) getResults(ctx context.Context, id uuid.UUID, coll *collector.Collector) (*[]result.Result, error) {
|
||||
func (s *Service) getResults(ctx context.Context, id uuid.UUID, coll *collector.Collector) ([]*result.Result, error) {
|
||||
docID := database.MustToDBUUID(id)
|
||||
|
||||
results, err := s.db.Queries.ListResultsByDocumentID(ctx, repository.ListResultsByDocumentIDParams{
|
||||
@@ -65,10 +65,10 @@ func (s *Service) getResults(ctx context.Context, id uuid.UUID, coll *collector.
|
||||
return nil, err
|
||||
}
|
||||
|
||||
cleanResults := make([]result.Result, len(results))
|
||||
cleanResults := make([]*result.Result, len(results))
|
||||
for index, dbResult := range results {
|
||||
cleanResults[index] = *result.Parse(&dbResult)
|
||||
cleanResults[index] = result.Parse(&dbResult)
|
||||
}
|
||||
|
||||
return &cleanResults, nil
|
||||
return cleanResults, nil
|
||||
}
|
||||
|
||||
@@ -8,10 +8,10 @@ type Result struct {
|
||||
value string
|
||||
}
|
||||
|
||||
func NewResult(value string) *Result {
|
||||
return &Result{value}
|
||||
func NewResult(value string) Result {
|
||||
return Result{value}
|
||||
}
|
||||
|
||||
func (r *Result) GetValue(ctx context.Context) (string, error) {
|
||||
func (r Result) GetValue(ctx context.Context) (string, error) {
|
||||
return r.value, nil
|
||||
}
|
||||
|
||||
@@ -24,12 +24,12 @@ func NewExtractor(db *database.Connection) Extractor {
|
||||
return Extractor{db}
|
||||
}
|
||||
|
||||
func (e Extractor) Process(ctx context.Context, query queryprocessor.Query, values *[]result.Value) (string, error) {
|
||||
if len(*values) != 1 {
|
||||
func (e Extractor) Process(ctx context.Context, query *queryprocessor.Query, values []result.Value) (string, error) {
|
||||
if len(values) != 1 {
|
||||
return "", fmt.Errorf("JSON Extraction requires 1 result")
|
||||
}
|
||||
|
||||
value, err := (*values)[0].GetValue(ctx)
|
||||
value, err := values[0].GetValue(ctx)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
@@ -102,10 +102,7 @@ func parseCreateQuery(q *queryprocessor.Create) (*createQuery, error) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
reqIDs := make([]pgtype.UUID, len(q.RequiredQueryIDs))
|
||||
for index, id := range q.RequiredQueryIDs {
|
||||
reqIDs[index] = database.MustToDBUUID(id)
|
||||
}
|
||||
reqIDs := database.MustToDBUUIDArray(q.RequiredQueryIDs)
|
||||
|
||||
return &createQuery{
|
||||
Type: t,
|
||||
|
||||
@@ -28,10 +28,7 @@ func (s *Service) Get(ctx context.Context, id uuid.UUID) (*Query, error) {
|
||||
}
|
||||
|
||||
func ParseDBGetQuery(q *repository.GetQueryRow) (*Query, error) {
|
||||
reqQueryIDs := make([]uuid.UUID, len(q.Requiredids))
|
||||
for index, id := range q.Requiredids {
|
||||
reqQueryIDs[index] = database.MustToUUID(id)
|
||||
}
|
||||
reqQueryIDs := database.MustToUUIDArray(q.Requiredids)
|
||||
|
||||
qType, err := queryprocessor.ParseDBType(q.Type)
|
||||
if err != nil {
|
||||
|
||||
+5
-10
@@ -5,15 +5,13 @@ import (
|
||||
"queryorchestration/internal/database"
|
||||
"queryorchestration/internal/database/repository"
|
||||
queryprocessor "queryorchestration/internal/queryProcessor"
|
||||
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
type ListFilters struct {
|
||||
Types []queryprocessor.Type
|
||||
}
|
||||
|
||||
func (s *Service) List(ctx context.Context, filters ListFilters) (*[]Query, error) {
|
||||
func (s *Service) List(ctx context.Context, filters ListFilters) ([]*Query, error) {
|
||||
// TODO - use filters
|
||||
|
||||
dbQueries, err := s.db.Queries.ListQueries(ctx)
|
||||
@@ -21,24 +19,21 @@ func (s *Service) List(ctx context.Context, filters ListFilters) (*[]Query, erro
|
||||
return nil, err
|
||||
}
|
||||
|
||||
queries := make([]Query, len(dbQueries))
|
||||
queries := make([]*Query, len(dbQueries))
|
||||
for index, query := range dbQueries {
|
||||
q, err := ParseDBListQuery(&query)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
queries[index] = *q
|
||||
queries[index] = q
|
||||
}
|
||||
|
||||
return &queries, nil
|
||||
return queries, nil
|
||||
}
|
||||
|
||||
func ParseDBListQuery(q *repository.ListQueriesRow) (*Query, error) {
|
||||
reqQueryIDs := make([]uuid.UUID, len(q.Requiredids))
|
||||
for index, id := range q.Requiredids {
|
||||
reqQueryIDs[index] = database.MustToUUID(id)
|
||||
}
|
||||
reqQueryIDs := database.MustToUUIDArray(q.Requiredids)
|
||||
|
||||
qType, err := queryprocessor.ParseDBType(q.Type)
|
||||
if err != nil {
|
||||
|
||||
@@ -4,8 +4,6 @@ import (
|
||||
"fmt"
|
||||
"queryorchestration/internal/database"
|
||||
"queryorchestration/internal/database/repository"
|
||||
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
func ParseDBNullType(qType repository.NullQuerytype) (Type, error) {
|
||||
@@ -52,10 +50,7 @@ func ToDBNullQueryType(t Type) (repository.NullQuerytype, error) {
|
||||
}
|
||||
|
||||
func ParseDBCollectorQuery(q *repository.GetCollectorQueriesRow) (*Query, error) {
|
||||
reqQueryIDs := make([]uuid.UUID, len(q.Requiredids))
|
||||
for index, id := range q.Requiredids {
|
||||
reqQueryIDs[index] = database.MustToUUID(id)
|
||||
}
|
||||
reqQueryIDs := database.MustToUUIDArray(q.Requiredids)
|
||||
|
||||
qType, err := ParseDBNullType(q.Type)
|
||||
if err != nil {
|
||||
|
||||
@@ -43,5 +43,5 @@ type Updator interface {
|
||||
}
|
||||
|
||||
type Processor interface {
|
||||
Process(ctx context.Context, query Query, values *[]result.Value) (string, error)
|
||||
Process(ctx context.Context, query *Query, values []result.Value) (string, error)
|
||||
}
|
||||
|
||||
@@ -7,9 +7,9 @@ import (
|
||||
)
|
||||
|
||||
func (q *Queue) getUnsyncedQueries() {
|
||||
for _, query := range *q.collectorQueries {
|
||||
for _, query := range q.collectorQueries {
|
||||
isSynced := false
|
||||
for _, result := range *q.results {
|
||||
for _, result := range q.results {
|
||||
if result.QueryID != query.ID || result.QueryVersion != query.Version {
|
||||
continue
|
||||
}
|
||||
@@ -20,7 +20,7 @@ func (q *Queue) getUnsyncedQueries() {
|
||||
continue
|
||||
}
|
||||
|
||||
q.Add(&query)
|
||||
q.Add(query)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -36,29 +36,29 @@ func (c *Queue) getCollectorQueries(ctx context.Context) error {
|
||||
return err
|
||||
}
|
||||
|
||||
cleanQueries := make([]queryprocessor.Query, len(queries))
|
||||
cleanQueries := make([]*queryprocessor.Query, len(queries))
|
||||
for index, dbQuery := range queries {
|
||||
cleanQuery, err := queryprocessor.ParseDBCollectorQuery(&dbQuery)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
cleanQueries[index] = *cleanQuery
|
||||
cleanQueries[index] = cleanQuery
|
||||
}
|
||||
|
||||
c.collectorQueries = &cleanQueries
|
||||
c.collectorQueries = cleanQueries
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (q *Queue) Add(qu *queryprocessor.Query) {
|
||||
dependentQueries := []queryprocessor.Query{}
|
||||
dependentQueries := []*queryprocessor.Query{}
|
||||
requiredIndex := -1
|
||||
|
||||
if q.unsyncedQueue == nil {
|
||||
q.unsyncedQueue = &[]queryprocessor.Query{}
|
||||
q.unsyncedQueue = []*queryprocessor.Query{}
|
||||
} else {
|
||||
for index, entry := range *q.unsyncedQueue {
|
||||
for index, entry := range q.unsyncedQueue {
|
||||
if entry.ID == qu.ID {
|
||||
return
|
||||
}
|
||||
@@ -71,7 +71,7 @@ func (q *Queue) Add(qu *queryprocessor.Query) {
|
||||
}
|
||||
}
|
||||
|
||||
for _, entry := range *q.collectorQueries {
|
||||
for _, entry := range q.collectorQueries {
|
||||
for _, id := range entry.RequiredQueryIDs {
|
||||
if qu.ID == id {
|
||||
dependentQueries = append(dependentQueries, entry)
|
||||
@@ -81,12 +81,12 @@ func (q *Queue) Add(qu *queryprocessor.Query) {
|
||||
}
|
||||
|
||||
if requiredIndex != -1 {
|
||||
*q.unsyncedQueue = append((*q.unsyncedQueue)[:requiredIndex+1], append([]queryprocessor.Query{*qu}, (*q.unsyncedQueue)[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...)
|
||||
q.unsyncedQueue = append([]*queryprocessor.Query{qu}, q.unsyncedQueue...)
|
||||
}
|
||||
|
||||
for _, entry := range dependentQueries {
|
||||
q.Add(&entry)
|
||||
q.Add(entry)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -18,7 +18,7 @@ func (q *Queue) Execute(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
for _, query := range *q.unsyncedQueue {
|
||||
for _, query := range q.unsyncedQueue {
|
||||
err := q.executeQuery(ctx, query)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -30,18 +30,18 @@ func (q *Queue) Execute(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (q *Queue) executeQuery(ctx context.Context, qu queryprocessor.Query) error {
|
||||
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 {
|
||||
for _, entry := range q.collectorQueries {
|
||||
if entry.ID == id {
|
||||
queryVersion = entry.Version
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
for _, entry := range *q.results {
|
||||
for _, entry := range q.results {
|
||||
if entry.QueryID == id && entry.QueryVersion == queryVersion {
|
||||
resultIDs[index] = database.MustToDBUUID(entry.ID)
|
||||
break
|
||||
@@ -64,7 +64,7 @@ func (q *Queue) executeQuery(ctx context.Context, qu queryprocessor.Query) error
|
||||
cleanValues[index] = cleanValue
|
||||
}
|
||||
|
||||
err = q.setResult(ctx, qu, &cleanValues)
|
||||
err = q.setResult(ctx, qu, cleanValues)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -74,7 +74,7 @@ func (q *Queue) executeQuery(ctx context.Context, qu queryprocessor.Query) error
|
||||
|
||||
func (q *Queue) getResultValue(res *repository.ListResultValuesByIDRow) (result.Value, error) {
|
||||
var queryType queryprocessor.Type
|
||||
for _, qu := range *q.collectorQueries {
|
||||
for _, qu := range q.collectorQueries {
|
||||
if qu.ID == database.MustToUUID(res.Queryid) {
|
||||
queryType = qu.Type
|
||||
}
|
||||
|
||||
@@ -9,7 +9,7 @@ import (
|
||||
"queryorchestration/internal/result"
|
||||
)
|
||||
|
||||
func (q *Queue) setResult(ctx context.Context, qu queryprocessor.Query, resultValues *[]result.Value) error {
|
||||
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
|
||||
@@ -32,7 +32,7 @@ func (q *Queue) setResult(ctx context.Context, qu queryprocessor.Query, resultVa
|
||||
return err
|
||||
}
|
||||
|
||||
*q.results = append(*q.results, result.Result{
|
||||
q.results = append(q.results, &result.Result{
|
||||
ID: id,
|
||||
QueryID: qu.ID,
|
||||
QueryVersion: qu.Version,
|
||||
|
||||
@@ -11,9 +11,9 @@ import (
|
||||
)
|
||||
|
||||
type Queue struct {
|
||||
unsyncedQueue *[]queryprocessor.Query
|
||||
collectorQueries *[]queryprocessor.Query
|
||||
results *[]result.Result
|
||||
unsyncedQueue []*queryprocessor.Query
|
||||
collectorQueries []*queryprocessor.Query
|
||||
results []*result.Result
|
||||
collector *collector.Collector
|
||||
db *database.Connection
|
||||
cleanVersion int32
|
||||
@@ -21,7 +21,7 @@ type Queue struct {
|
||||
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) {
|
||||
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,
|
||||
@@ -41,6 +41,6 @@ func New(ctx context.Context, db *database.Connection, coll *collector.Collector
|
||||
return &queue, nil
|
||||
}
|
||||
|
||||
func (q *Queue) GetQueue() []queryprocessor.Query {
|
||||
return *q.unsyncedQueue
|
||||
func (q *Queue) GetQueue() []*queryprocessor.Query {
|
||||
return q.unsyncedQueue
|
||||
}
|
||||
|
||||
@@ -16,20 +16,20 @@ func TestContextFull(t *testing.T) {
|
||||
|
||||
extractor := contextfull.NewExtractor()
|
||||
|
||||
query := queryprocessor.Query{
|
||||
query := &queryprocessor.Query{
|
||||
ID: uuid.New(),
|
||||
Type: queryprocessor.TypeJsonExtractor,
|
||||
RequiredQueryIDs: []uuid.UUID{},
|
||||
Version: int32(1),
|
||||
}
|
||||
|
||||
values := &[]result.Value{}
|
||||
values := []result.Value{}
|
||||
|
||||
value, err := extractor.Process(ctx, query, values)
|
||||
assert.Nil(t, err)
|
||||
assert.Equal(t, "", value)
|
||||
|
||||
values = &[]result.Value{
|
||||
values = []result.Value{
|
||||
contextfull.NewResult("example_result"),
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,49 @@
|
||||
package database_test
|
||||
|
||||
import (
|
||||
"queryorchestration/internal/database"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestMustToDBUUID(t *testing.T) {
|
||||
id := uuid.New()
|
||||
|
||||
dbID := database.MustToDBUUID(id)
|
||||
|
||||
assert.Equal(t, true, dbID.Valid)
|
||||
assert.Equal(t, id.String(), uuid.UUID(dbID.Bytes).String())
|
||||
}
|
||||
|
||||
func TestMustToDBUUIDArray(t *testing.T) {
|
||||
ids := []uuid.UUID{uuid.New(), uuid.New()}
|
||||
|
||||
dbIDs := database.MustToDBUUIDArray(ids)
|
||||
|
||||
assert.Equal(t, len(ids), len(dbIDs))
|
||||
for index, id := range dbIDs {
|
||||
assert.Equal(t, database.MustToDBUUID(ids[index]), id)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMustToUUID(t *testing.T) {
|
||||
dbID := database.MustToDBUUID(uuid.New())
|
||||
|
||||
id := database.MustToUUID(dbID)
|
||||
|
||||
assert.Equal(t, id.String(), uuid.UUID(dbID.Bytes).String())
|
||||
}
|
||||
|
||||
func TestMustToUUIDArray(t *testing.T) {
|
||||
dbIDs := []pgtype.UUID{database.MustToDBUUID(uuid.New()), database.MustToDBUUID(uuid.New())}
|
||||
|
||||
ids := database.MustToUUIDArray(dbIDs)
|
||||
|
||||
assert.Equal(t, len(ids), len(dbIDs))
|
||||
for index, id := range dbIDs {
|
||||
assert.Equal(t, database.MustToDBUUID(ids[index]), id)
|
||||
}
|
||||
}
|
||||
@@ -1,18 +0,0 @@
|
||||
package database_test
|
||||
|
||||
import (
|
||||
"queryorchestration/internal/database"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func TestMustToDBUUID(t *testing.T) {
|
||||
id := uuid.New()
|
||||
|
||||
dbID := database.MustToDBUUID(id)
|
||||
|
||||
assert.Equal(t, true, dbID.Valid)
|
||||
assert.Equal(t, id.String(), uuid.UUID(dbID.Bytes).String())
|
||||
}
|
||||
@@ -32,7 +32,7 @@ func TestJSONProcess(t *testing.T) {
|
||||
|
||||
extractor := jsonextractor.NewExtractor(db)
|
||||
|
||||
query := queryprocessor.Query{
|
||||
query := &queryprocessor.Query{
|
||||
ID: uuid.New(),
|
||||
Type: queryprocessor.TypeJsonExtractor,
|
||||
RequiredQueryIDs: []uuid.UUID{},
|
||||
@@ -41,7 +41,7 @@ func TestJSONProcess(t *testing.T) {
|
||||
entryValue := "value"
|
||||
|
||||
jsonString := fmt.Sprintf("{\"key\": \"%s\"}", entryValue)
|
||||
values := &[]result.Value{
|
||||
values := []result.Value{
|
||||
contextfull.NewResult(jsonString),
|
||||
}
|
||||
|
||||
@@ -59,7 +59,7 @@ func TestJSONProcess(t *testing.T) {
|
||||
|
||||
entryValue = ""
|
||||
jsonString = fmt.Sprintf("{\"key\": \"%s\"}", entryValue)
|
||||
values = &[]result.Value{
|
||||
values = []result.Value{
|
||||
contextfull.NewResult(jsonString),
|
||||
}
|
||||
pool.ExpectQuery("name: GetQueryConfig :one").WithArgs(database.MustToDBUUID(query.ID), query.Version).
|
||||
@@ -74,7 +74,7 @@ func TestJSONProcess(t *testing.T) {
|
||||
|
||||
entryValue = "1"
|
||||
jsonString = fmt.Sprintf("{\"key\": %s", entryValue)
|
||||
values = &[]result.Value{
|
||||
values = []result.Value{
|
||||
contextfull.NewResult(jsonString),
|
||||
}
|
||||
pool.ExpectQuery("name: GetQueryConfig :one").WithArgs(database.MustToDBUUID(query.ID), query.Version).
|
||||
@@ -103,7 +103,7 @@ func TestJSONProcessJSON(t *testing.T) {
|
||||
|
||||
extractor := jsonextractor.NewExtractor(db)
|
||||
|
||||
query := queryprocessor.Query{
|
||||
query := &queryprocessor.Query{
|
||||
ID: uuid.New(),
|
||||
Type: queryprocessor.TypeJsonExtractor,
|
||||
RequiredQueryIDs: []uuid.UUID{},
|
||||
@@ -112,7 +112,7 @@ func TestJSONProcessJSON(t *testing.T) {
|
||||
entryValue := "value"
|
||||
|
||||
jsonString := fmt.Sprintf("{\"key\": \"%s\"}", entryValue)
|
||||
values := &[]result.Value{
|
||||
values := []result.Value{
|
||||
contextfull.NewResult(jsonString),
|
||||
}
|
||||
|
||||
@@ -191,19 +191,19 @@ func TestJSONProcessResults(t *testing.T) {
|
||||
|
||||
extractor := jsonextractor.NewExtractor(db)
|
||||
|
||||
query := queryprocessor.Query{
|
||||
query := &queryprocessor.Query{
|
||||
ID: uuid.New(),
|
||||
Type: queryprocessor.TypeJsonExtractor,
|
||||
RequiredQueryIDs: []uuid.UUID{},
|
||||
Version: int32(1),
|
||||
}
|
||||
|
||||
results := &[]result.Value{}
|
||||
results := []result.Value{}
|
||||
value, err := extractor.Process(ctx, query, results)
|
||||
assert.EqualError(t, err, "JSON Extraction requires 1 result")
|
||||
assert.Empty(t, value)
|
||||
|
||||
results = &[]result.Value{
|
||||
results = []result.Value{
|
||||
contextfull.NewResult(""),
|
||||
contextfull.NewResult(""),
|
||||
}
|
||||
@@ -211,7 +211,7 @@ func TestJSONProcessResults(t *testing.T) {
|
||||
assert.EqualError(t, err, "JSON Extraction requires 1 result")
|
||||
assert.Empty(t, value)
|
||||
|
||||
results = &[]result.Value{
|
||||
results = []result.Value{
|
||||
contextfull.NewResult(""),
|
||||
contextfull.NewResult(""),
|
||||
contextfull.NewResult(""),
|
||||
|
||||
@@ -10,7 +10,6 @@ import (
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
"github.com/pashagolub/pgxmock/v3"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
@@ -46,10 +45,6 @@ func TestCreate(t *testing.T) {
|
||||
Config: q.Config,
|
||||
}
|
||||
|
||||
dbReqIDs := make([]pgtype.UUID, len(q.RequiredQueryIDs))
|
||||
for index, id := range q.RequiredQueryIDs {
|
||||
dbReqIDs[index] = database.MustToDBUUID(id)
|
||||
}
|
||||
dbType, err := queryprocessor.ToDBQueryType(create.Type)
|
||||
assert.Nil(t, err)
|
||||
|
||||
|
||||
@@ -29,7 +29,7 @@ func TestList(t *testing.T) {
|
||||
svc := query.New(db)
|
||||
|
||||
config := "{\"path\":\"example_path\"}"
|
||||
q := query.Query{
|
||||
q := &query.Query{
|
||||
ID: uuid.New(),
|
||||
Type: queryprocessor.TypeJsonExtractor,
|
||||
ActiveVersion: int32(1),
|
||||
@@ -55,5 +55,5 @@ func TestList(t *testing.T) {
|
||||
resList, err := svc.List(ctx, filters)
|
||||
assert.Nil(t, err)
|
||||
|
||||
assert.EqualExportedValues(t, []query.Query{q}, *resList)
|
||||
assert.EqualExportedValues(t, []*query.Query{q}, resList)
|
||||
}
|
||||
|
||||
@@ -9,7 +9,6 @@ import (
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
"github.com/pashagolub/pgxmock/v3"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
@@ -47,10 +46,7 @@ func TestUpdate(t *testing.T) {
|
||||
Config: config,
|
||||
}
|
||||
|
||||
dbReqIDs := make([]pgtype.UUID, len(existing.RequiredQueryIDs))
|
||||
for index, id := range existing.RequiredQueryIDs {
|
||||
dbReqIDs[index] = database.MustToDBUUID(id)
|
||||
}
|
||||
dbReqIDs := database.MustToDBUUIDArray(existing.RequiredQueryIDs)
|
||||
|
||||
pool.ExpectQuery("name: GetQuery :one").WithArgs(database.MustToDBUUID(update.ID)).WillReturnRows(
|
||||
pgxmock.NewRows([]string{"id", "type", "activeVersion", "latestVersion", "config", "requiredIds"}).
|
||||
|
||||
@@ -84,14 +84,14 @@ func TestQueue(t *testing.T) {
|
||||
pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID).WillReturnRows(rows)
|
||||
|
||||
contextResultID := uuid.New()
|
||||
results := []result.Result{
|
||||
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{
|
||||
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},
|
||||
@@ -103,7 +103,7 @@ func TestQueue(t *testing.T) {
|
||||
cleanVersion := int32(1)
|
||||
textVersion := int32(1)
|
||||
|
||||
q, err := queryQueue.New(ctx, db, coll, &results, docID, cleanVersion, textVersion)
|
||||
q, err := queryQueue.New(ctx, db, coll, results, docID, cleanVersion, textVersion)
|
||||
assert.Nil(t, err)
|
||||
assert.Equal(t, expectedQueries, q.GetQueue())
|
||||
|
||||
@@ -212,12 +212,12 @@ func TestQueueFail(t *testing.T) {
|
||||
pool.ExpectQuery("name: GetCollectorQueries :many").WithArgs(dbCollectorID).
|
||||
WillReturnError(errors.New(errr))
|
||||
|
||||
results := []result.Result{}
|
||||
results := []*result.Result{}
|
||||
|
||||
docID := uuid.New()
|
||||
cleanVersion := int32(1)
|
||||
textVersion := int32(1)
|
||||
|
||||
_, err = queryQueue.New(ctx, db, coll, &results, docID, cleanVersion, textVersion)
|
||||
_, err = queryQueue.New(ctx, db, coll, results, docID, cleanVersion, textVersion)
|
||||
assert.EqualError(t, err, errr)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user