5c253b3592
Feature/docinit * createunittests * cleanup * skip * cleanup
280 lines
6.4 KiB
Go
280 lines
6.4 KiB
Go
package collector
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"queryorchestration/internal/database"
|
|
"queryorchestration/internal/database/repository"
|
|
"queryorchestration/internal/validation"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/jackc/pgx/v5/pgtype"
|
|
)
|
|
|
|
type UpdateParams struct {
|
|
JobID uuid.UUID
|
|
ActiveVersion *int32
|
|
MinCleanVersion *int32
|
|
MinTextVersion *int32
|
|
Fields *map[string]uuid.UUID
|
|
}
|
|
|
|
func (s *Service) UpdateByJobId(ctx context.Context, params *UpdateParams) error {
|
|
current, err := s.GetByJobID(ctx, params.JobID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
dbparams, err := s.getUpdateParams(ctx, current, params)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = s.submitUpdate(ctx, current, dbparams)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
type dbUpdateParams struct {
|
|
JobID pgtype.UUID
|
|
ActiveVersion *int32
|
|
MinCleanVersion *int32
|
|
MinTextVersion *int32
|
|
Fields *map[string]pgtype.UUID
|
|
}
|
|
|
|
func (s *Service) getUpdateParams(ctx context.Context, current *Collector, params *UpdateParams) (*dbUpdateParams, error) {
|
|
err := s.normalizeCodeVersions(current, params)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
fs, err := s.normalizeFieldsToDB(ctx, params.Fields)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
err = s.normalizeActiveVersion(current, params)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if params.ActiveVersion == nil &&
|
|
params.MinCleanVersion == nil &&
|
|
params.MinTextVersion == nil &&
|
|
fs == nil {
|
|
return nil, errors.New("no changes")
|
|
}
|
|
|
|
return &dbUpdateParams{
|
|
JobID: database.MustToDBUUID(params.JobID),
|
|
ActiveVersion: params.ActiveVersion,
|
|
MinCleanVersion: params.MinCleanVersion,
|
|
MinTextVersion: params.MinTextVersion,
|
|
Fields: fs,
|
|
}, nil
|
|
}
|
|
|
|
func (s *Service) normalizeCodeVersions(current *Collector, params *UpdateParams) error {
|
|
if current == nil {
|
|
return errors.New("current collector required")
|
|
}
|
|
|
|
if params == nil || params.MinCleanVersion == nil && params.MinTextVersion == nil {
|
|
return nil
|
|
}
|
|
|
|
if (params.MinCleanVersion == nil || *params.MinCleanVersion == current.MinCleanVersion) &&
|
|
(params.MinTextVersion == nil || *params.MinTextVersion == current.MinTextVersion) {
|
|
params.MinCleanVersion = nil
|
|
params.MinTextVersion = nil
|
|
return nil
|
|
}
|
|
|
|
if params.MinCleanVersion == nil {
|
|
params.MinCleanVersion = ¤t.MinCleanVersion
|
|
} else {
|
|
err := s.svc.Clean.IsValidVersion(*params.MinCleanVersion)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
if params.MinTextVersion == nil {
|
|
params.MinTextVersion = ¤t.MinTextVersion
|
|
} else {
|
|
err := s.svc.Text.IsValidVersion(*params.MinTextVersion)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *Service) normalizeActiveVersion(current *Collector, params *UpdateParams) error {
|
|
if current == nil {
|
|
return errors.New("current collector required")
|
|
}
|
|
|
|
if params == nil || params.ActiveVersion == nil {
|
|
return nil
|
|
}
|
|
|
|
err := validation.NormalizeInClosedInterval(¶ms.ActiveVersion, current.ActiveVersion, 1, current.LatestVersion+1)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *Service) normalizeFieldsToDB(ctx context.Context, ofields *map[string]uuid.UUID) (*map[string]pgtype.UUID, error) {
|
|
if ofields == nil || *ofields == nil {
|
|
return nil, nil
|
|
} else if len(*ofields) == 0 {
|
|
return nil, nil
|
|
}
|
|
|
|
dbm := map[string]pgtype.UUID{}
|
|
for name, id := range *ofields {
|
|
dbm[name] = database.MustToDBUUID(id)
|
|
}
|
|
|
|
dbids := []pgtype.UUID{}
|
|
for _, id := range dbm {
|
|
dbids = append(dbids, id)
|
|
}
|
|
|
|
dedup := validation.DeduplicateArray(dbids)
|
|
|
|
if len(dedup) != len(dbids) {
|
|
return nil, errors.New("duplicate output fields")
|
|
}
|
|
|
|
exist, err := s.cfg.GetDBQueries().AllQueriesExist(ctx, dbids)
|
|
if err != nil {
|
|
return nil, err
|
|
} else if !exist {
|
|
return nil, errors.New("not all required ids are present")
|
|
}
|
|
|
|
return &dbm, nil
|
|
}
|
|
|
|
func (s *Service) submitUpdate(ctx context.Context, current *Collector, params *dbUpdateParams) error {
|
|
err := s.cfg.ExecuteDBTransaction(ctx, func(ctx context.Context, qtx *repository.Queries) error {
|
|
latestVersion := current.LatestVersion + 1
|
|
id := database.MustToDBUUID(current.ID)
|
|
|
|
if params.MinCleanVersion != nil || params.MinTextVersion != nil {
|
|
err := qtx.RemoveCollectorCodeVersion(ctx, &repository.RemoveCollectorCodeVersionParams{
|
|
Collectorid: id,
|
|
Removedversion: &latestVersion,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
err = qtx.AddCollectorCodeVersion(ctx, &repository.AddCollectorCodeVersionParams{
|
|
Collectorid: id,
|
|
Mincleanversion: *params.MinCleanVersion,
|
|
Mintextversion: *params.MinTextVersion,
|
|
Addedversion: latestVersion,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
removeIDs := getRemoveFields(current.Fields, params.Fields)
|
|
for _, field := range removeIDs {
|
|
err := qtx.RemoveCollectorQuery(ctx, &repository.RemoveCollectorQueryParams{
|
|
Collectorid: id,
|
|
Queryid: field,
|
|
Removedversion: &latestVersion,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
addIDs := getAddFields(current.Fields, params.Fields)
|
|
for key, field := range addIDs {
|
|
err := qtx.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{
|
|
Collectorid: id,
|
|
Name: key,
|
|
Queryid: field,
|
|
Addedversion: latestVersion,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
activeVersion := params.ActiveVersion
|
|
if activeVersion == nil {
|
|
activeVersion = ¤t.ActiveVersion
|
|
}
|
|
|
|
err := qtx.UpdateCollector(ctx, &repository.UpdateCollectorParams{
|
|
ID: id,
|
|
Latestversion: latestVersion,
|
|
Activeversion: *activeVersion,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func getRemoveFields(current map[string]uuid.UUID, update *map[string]pgtype.UUID) []pgtype.UUID {
|
|
diff := []pgtype.UUID{}
|
|
if update == nil {
|
|
return diff
|
|
}
|
|
|
|
for ckey, cid := range current {
|
|
found := false
|
|
for ukey, uid := range *update {
|
|
if cid == database.MustToUUID(uid) &&
|
|
ckey == ukey {
|
|
found = true
|
|
}
|
|
}
|
|
if !found {
|
|
diff = append(diff, database.MustToDBUUID(cid))
|
|
}
|
|
}
|
|
|
|
return diff
|
|
}
|
|
|
|
func getAddFields(current map[string]uuid.UUID, update *map[string]pgtype.UUID) map[string]pgtype.UUID {
|
|
diff := map[string]pgtype.UUID{}
|
|
if update == nil {
|
|
return diff
|
|
}
|
|
|
|
for ckey, cid := range current {
|
|
for ukey, uid := range *update {
|
|
if ckey == ukey && cid != database.MustToUUID(uid) ||
|
|
current[ukey] == uuid.Nil {
|
|
diff[ukey] = uid
|
|
}
|
|
}
|
|
}
|
|
|
|
return diff
|
|
}
|