collectscollector
This commit is contained in:
@@ -11,35 +11,56 @@ import (
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
)
|
||||
|
||||
const getQueriesFromJobID = `-- name: GetQueriesFromJobID :many
|
||||
SELECT id, type, requiredQueryId, queryVersion, minCleanVersion, minTextVersion FROM collectorQueryDependencyTree WHERE jobId = $1
|
||||
const getCollectorFromJobID = `-- name: GetCollectorFromJobID :one
|
||||
SELECT id, jobId, minCleanVersion, minTextVersion FROM collectors WHERE jobId = $1 LIMIT 1
|
||||
`
|
||||
|
||||
type GetQueriesFromJobIDRow struct {
|
||||
type GetCollectorFromJobIDRow struct {
|
||||
ID pgtype.UUID
|
||||
Type Querytype
|
||||
Requiredqueryid pgtype.UUID
|
||||
Queryversion int32
|
||||
Jobid pgtype.UUID
|
||||
Mincleanversion int32
|
||||
Mintextversion int32
|
||||
}
|
||||
|
||||
func (q *Queries) GetQueriesFromJobID(ctx context.Context, jobid pgtype.UUID) ([]GetQueriesFromJobIDRow, error) {
|
||||
rows, err := q.db.Query(ctx, getQueriesFromJobID, jobid)
|
||||
func (q *Queries) GetCollectorFromJobID(ctx context.Context, jobid pgtype.UUID) (GetCollectorFromJobIDRow, error) {
|
||||
row := q.db.QueryRow(ctx, getCollectorFromJobID, jobid)
|
||||
var i GetCollectorFromJobIDRow
|
||||
err := row.Scan(
|
||||
&i.ID,
|
||||
&i.Jobid,
|
||||
&i.Mincleanversion,
|
||||
&i.Mintextversion,
|
||||
)
|
||||
return i, err
|
||||
}
|
||||
|
||||
const getCollectorQueries = `-- name: GetCollectorQueries :many
|
||||
SELECT collectorId, queryId, type, requiredQueryId, queryVersion FROM collectorQueryDependencyTree WHERE collectorId = $1
|
||||
`
|
||||
|
||||
type GetCollectorQueriesRow struct {
|
||||
Collectorid pgtype.UUID
|
||||
Queryid pgtype.UUID
|
||||
Type NullQuerytype
|
||||
Requiredqueryid pgtype.UUID
|
||||
Queryversion pgtype.Int4
|
||||
}
|
||||
|
||||
func (q *Queries) GetCollectorQueries(ctx context.Context, collectorid pgtype.UUID) ([]GetCollectorQueriesRow, error) {
|
||||
rows, err := q.db.Query(ctx, getCollectorQueries, collectorid)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
var items []GetQueriesFromJobIDRow
|
||||
var items []GetCollectorQueriesRow
|
||||
for rows.Next() {
|
||||
var i GetQueriesFromJobIDRow
|
||||
var i GetCollectorQueriesRow
|
||||
if err := rows.Scan(
|
||||
&i.ID,
|
||||
&i.Collectorid,
|
||||
&i.Queryid,
|
||||
&i.Type,
|
||||
&i.Requiredqueryid,
|
||||
&i.Queryversion,
|
||||
&i.Mincleanversion,
|
||||
&i.Mintextversion,
|
||||
); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -53,12 +53,9 @@ func (ns NullQuerytype) Value() (driver.Value, error) {
|
||||
}
|
||||
|
||||
type Activecollectorquery struct {
|
||||
ID pgtype.UUID
|
||||
Jobid pgtype.UUID
|
||||
Mincleanversion int32
|
||||
Mintextversion int32
|
||||
Activeversion int32
|
||||
Queryid pgtype.UUID
|
||||
Collectorid pgtype.UUID
|
||||
Activeversion int32
|
||||
Queryid pgtype.UUID
|
||||
}
|
||||
|
||||
type Activequeryrequirement struct {
|
||||
@@ -87,13 +84,11 @@ type Collectorquery struct {
|
||||
}
|
||||
|
||||
type Collectorquerydependencytree struct {
|
||||
Collectorid pgtype.UUID
|
||||
ID pgtype.UUID
|
||||
Jobid pgtype.UUID
|
||||
Type Querytype
|
||||
Requiredqueryid pgtype.UUID
|
||||
Queryversion int32
|
||||
Mincleanversion int32
|
||||
Mintextversion int32
|
||||
}
|
||||
|
||||
type Query struct {
|
||||
|
||||
@@ -12,9 +12,15 @@ import (
|
||||
)
|
||||
|
||||
const listResultsByDocumentID = `-- name: ListResultsByDocumentID :many
|
||||
SELECT id, queryId, cleanVersion, textVersion, queryVersion FROM results where documentId = $1
|
||||
SELECT id, queryId, cleanVersion, textVersion, queryVersion FROM results where documentId = $1 and cleanVersion >= $2 and textVersion >= $3
|
||||
`
|
||||
|
||||
type ListResultsByDocumentIDParams struct {
|
||||
Documentid pgtype.UUID
|
||||
Cleanversion int32
|
||||
Textversion int32
|
||||
}
|
||||
|
||||
type ListResultsByDocumentIDRow struct {
|
||||
ID pgtype.UUID
|
||||
Queryid pgtype.UUID
|
||||
@@ -23,8 +29,8 @@ type ListResultsByDocumentIDRow struct {
|
||||
Queryversion int32
|
||||
}
|
||||
|
||||
func (q *Queries) ListResultsByDocumentID(ctx context.Context, documentid pgtype.UUID) ([]ListResultsByDocumentIDRow, error) {
|
||||
rows, err := q.db.Query(ctx, listResultsByDocumentID, documentid)
|
||||
func (q *Queries) ListResultsByDocumentID(ctx context.Context, arg ListResultsByDocumentIDParams) ([]ListResultsByDocumentIDRow, error) {
|
||||
rows, err := q.db.Query(ctx, listResultsByDocumentID, arg.Documentid, arg.Cleanversion, arg.Textversion)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -19,6 +19,19 @@ type Document struct {
|
||||
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,
|
||||
@@ -26,46 +39,78 @@ func New(ctx context.Context, db *repository.Queries) *Service {
|
||||
}
|
||||
|
||||
func (s *Service) Sync(ctx context.Context, doc Document) error {
|
||||
jobID, err := database.ToDBUUID(doc.JobID)
|
||||
queries, err := s.getUnsyncedOrderedQueries(ctx, doc)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
queries, err := s.db.GetQueriesFromJobID(ctx, jobID)
|
||||
err = s.executeQueries(queries)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
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)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// TODO grab only the results with the latest clean and text version
|
||||
results, err := s.db.ListResultsByDocumentID(ctx, docID, minCleanVersion, minTextVersion)
|
||||
results, err := s.db.ListResultsByDocumentID(ctx, repository.ListResultsByDocumentIDParams{
|
||||
Documentid: docID,
|
||||
Textversion: collector.Mintextversion,
|
||||
Cleanversion: collector.Mincleanversion,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
|
||||
queue := []query{}
|
||||
|
||||
for _, query := range queries {
|
||||
isSynced := false
|
||||
for _, result := range results {
|
||||
if result.Queryid != query.ID && result.Queryversion != query.Queryversion {
|
||||
if result.Queryid != query.Queryid || int(result.Queryversion) != int(query.Queryversion.Int32) {
|
||||
continue
|
||||
}
|
||||
if result.Cleanversion > query.Mincleanversion && result.Textversion > query.Mintextversion {
|
||||
isSynced = true
|
||||
}
|
||||
isSynced = true
|
||||
break
|
||||
}
|
||||
if isSynced {
|
||||
continue
|
||||
}
|
||||
log.Print("TODO sync me pls")
|
||||
// Check if exists and synced
|
||||
// IF NOT synced then add to queue, as well as all dependents recursively
|
||||
// TODO add query to beginning, add all dependents until one that is already in the queue
|
||||
}
|
||||
|
||||
return queue, 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
|
||||
}
|
||||
// TODO what to do with results that are not used
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user