queuecreation

This commit is contained in:
Michael McGuinness
2024-12-20 17:35:33 +00:00
parent 6b07844d46
commit f95114e11a
11 changed files with 603 additions and 89 deletions
+47
View File
@@ -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
}
+1 -1
View File
@@ -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
+59 -10
View File
@@ -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
}
+15
View File
@@ -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[:]))
}
+36 -68
View File
@@ -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
}
+38
View File
@@ -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}
}
+251
View File
@@ -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
}