diff --git a/internal/jsonExtractor/service.go b/internal/jsonExtractor/service.go new file mode 100644 index 00000000..39b37bbc --- /dev/null +++ b/internal/jsonExtractor/service.go @@ -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 +} diff --git a/internal/query/create.go b/internal/query/create.go new file mode 100644 index 00000000..9fe49136 --- /dev/null +++ b/internal/query/create.go @@ -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) + } +} diff --git a/internal/query/database.go b/internal/query/database.go index 73262fef..18758087 100644 --- a/internal/query/database.go +++ b/internal/query/database.go @@ -29,7 +29,6 @@ func ToDBQueryType(t QueryType) repository.NullQuerytype { switch t { case QueryTypeJsonExtractor: dbType = repository.QuerytypeJsonExtractor - break default: dbType = repository.QuerytypeJsonExtractor } diff --git a/internal/query/execute.go b/internal/query/execute.go new file mode 100644 index 00000000..652fdd23 --- /dev/null +++ b/internal/query/execute.go @@ -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 +} diff --git a/internal/query/result.go b/internal/query/result.go new file mode 100644 index 00000000..44a16833 --- /dev/null +++ b/internal/query/result.go @@ -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") + } +} diff --git a/internal/query/service.go b/internal/query/service.go index f055b990..e9808e68 100644 --- a/internal/query/service.go +++ b/internal/query/service.go @@ -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 -} diff --git a/test/unit/internal/document/sync_test.go b/test/unit/internal/document/sync_test.go index 33213e43..6ecfd9f4 100644 --- a/test/unit/internal/document/sync_test.go +++ b/test/unit/internal/document/sync_test.go @@ -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) }