firstroundsetupforquerytypes
This commit is contained in:
@@ -0,0 +1,9 @@
|
||||
package jsonextractor
|
||||
|
||||
import "github.com/google/uuid"
|
||||
|
||||
func Process(queryId uuid.UUID, queryVersion int32, jsonString string) (string, error) {
|
||||
// TODO - get full query information
|
||||
// TODO - process input
|
||||
return "", nil
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
package query
|
||||
|
||||
import (
|
||||
"context"
|
||||
"gotemplate/internal/database"
|
||||
)
|
||||
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -29,7 +29,6 @@ func ToDBQueryType(t QueryType) repository.NullQuerytype {
|
||||
switch t {
|
||||
case QueryTypeJsonExtractor:
|
||||
dbType = repository.QuerytypeJsonExtractor
|
||||
break
|
||||
default:
|
||||
dbType = repository.QuerytypeJsonExtractor
|
||||
}
|
||||
|
||||
@@ -0,0 +1,110 @@
|
||||
package query
|
||||
|
||||
import (
|
||||
"context"
|
||||
"gotemplate/internal/database"
|
||||
"gotemplate/internal/database/repository"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"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, 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 {
|
||||
value, err := q.prepareResult(ctx, query, resultValues)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package query
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
jsonextractor "gotemplate/internal/jsonExtractor"
|
||||
)
|
||||
|
||||
func (q *Queue) prepareResult(ctx context.Context, query Query, resultValues *[]ResultValue) (string, error) {
|
||||
switch query.Type {
|
||||
case QueryTypeJsonExtractor:
|
||||
return jsonextractor.Process(query.ID, query.Version, (*resultValues)[0].Value)
|
||||
default:
|
||||
return "", fmt.Errorf("attempting to process invalid query type")
|
||||
}
|
||||
}
|
||||
@@ -3,11 +3,9 @@ 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
|
||||
@@ -69,183 +67,3 @@ func NewQueue(ctx context.Context, db *repository.Queries, coll *collector.Colle
|
||||
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
|
||||
}
|
||||
|
||||
@@ -5,7 +5,6 @@ import (
|
||||
"gotemplate/internal/database"
|
||||
"gotemplate/internal/database/repository"
|
||||
"gotemplate/internal/document"
|
||||
"gotemplate/internal/sync"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
@@ -25,7 +24,6 @@ func TestSync(t *testing.T) {
|
||||
|
||||
queries := repository.New(db)
|
||||
docSvc := document.New(ctx, queries)
|
||||
svc := sync.New(ctx, queries, docSvc)
|
||||
|
||||
doc := document.Document{
|
||||
ID: uuid.New(),
|
||||
@@ -59,7 +57,7 @@ func TestSync(t *testing.T) {
|
||||
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 = docSvc.Sync(ctx, &doc)
|
||||
|
||||
assert.Nil(t, err)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user