diff --git a/api/queryService/job_test.go b/api/queryService/job_test.go index 5b8475b5..326416d6 100644 --- a/api/queryService/job_test.go +++ b/api/queryService/job_test.go @@ -57,6 +57,7 @@ func TestCreateJob(t *testing.T) { ctx := e.NewContext(req, rec) id := uuid.New() + collId := uuid.New() pool.ExpectQuery("name: CreateJob :one").WithArgs(database.MustToDBUUID(body.ClientId)).WillReturnRows( pgxmock.NewRows([]string{"id"}). @@ -65,8 +66,12 @@ func TestCreateJob(t *testing.T) { pool.ExpectBeginTx(pgx.TxOptions{}) pool.ExpectQuery("name: CreateCollector :one").WithArgs(database.MustToDBUUID(id)).WillReturnRows( pgxmock.NewRows([]string{"id"}). - AddRow(database.MustToDBUUID(uuid.New())), + AddRow(database.MustToDBUUID(collId)), ) + pool.ExpectExec("name: AddLatestCollectorVersion :exec").WithArgs(database.MustToDBUUID(collId)). + WillReturnResult(pgxmock.NewResult("", 1)) + pool.ExpectExec("name: AddActiveCollectorVersion :exec").WithArgs(database.MustToDBUUID(collId), int32(1)). + WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() err = cons.CreateJob(ctx) diff --git a/api/queryService/jobcollector_test.go b/api/queryService/jobcollector_test.go index 8057d215..ba287d76 100644 --- a/api/queryService/jobcollector_test.go +++ b/api/queryService/jobcollector_test.go @@ -94,14 +94,16 @@ func TestUpdateJobCollector(t *testing.T) { ) pool.ExpectBeginTx(pgx.TxOptions{}) rv := current.LatestVersion + 1 + pool.ExpectExec("name: AddLatestCollectorVersion :exec").WithArgs(database.MustToDBUUID(current.ID)). + WillReturnResult(pgxmock.NewResult("", 1)) + pool.ExpectExec("name: AddActiveCollectorVersion :exec").WithArgs(database.MustToDBUUID(current.ID), int32(2)). + WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: RemoveCollectorCodeVersion :exec").WithArgs(&rv, database.MustToDBUUID(current.ID)). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: AddCollectorCodeVersion :exec").WithArgs(database.MustToDBUUID(current.ID), int32(5), *body.MinimumCleanerVersion, int32(0)). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(current.ID), "a", database.MustToDBUUID((*body.Fields)[0].QueryId), int32(5)). WillReturnResult(pgxmock.NewResult("", 1)) - pool.ExpectExec("name: UpdateCollector :exec").WithArgs(int32(5), int32(2), database.MustToDBUUID(current.ID)). - WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() mockSQS.EXPECT(). diff --git a/database/migrations/00000000000002_queries.up.sql b/database/migrations/00000000000002_queries.up.sql index 834d3017..97c32c38 100644 --- a/database/migrations/00000000000002_queries.up.sql +++ b/database/migrations/00000000000002_queries.up.sql @@ -25,7 +25,7 @@ BEGIN END; $$ LANGUAGE plpgsql; -CREATE TRIGGER setVersionNumberTrigger +CREATE TRIGGER setQueryVersionNumberTrigger BEFORE INSERT ON queryVersions FOR EACH ROW EXECUTE FUNCTION setQueryVersionNumber(); diff --git a/database/migrations/00000000000005_collectors.up.sql b/database/migrations/00000000000005_collectors.up.sql index 7b5c6004..dd8b0555 100644 --- a/database/migrations/00000000000005_collectors.up.sql +++ b/database/migrations/00000000000005_collectors.up.sql @@ -1,12 +1,42 @@ CREATE TABLE collectors ( id uuid primary key DEFAULT uuid_generate_v7(), jobId uuid not null, - latestVersion int not null default 1, - activeVersion int not null default 1, foreign key (jobId) references jobs(id), unique (jobId) ); +CREATE TABLE collectorVersions ( + collectorId uuid not null, + id int not null, + addedAt timestamp not null default current_timestamp, + primary key (id, collectorId), + foreign key (collectorId) references collectors(id) +); + +CREATE OR REPLACE FUNCTION setCollectorVersionNumber() +RETURNS TRIGGER AS $$ +BEGIN + SELECT COALESCE(MAX(id), 0) + 1 + INTO NEW.id + FROM collectorVersions + WHERE collectorId = NEW.collectorId; + + RETURN NEW; +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER setCollectorVersionNumberTrigger +BEFORE INSERT ON collectorVersions +FOR EACH ROW +EXECUTE FUNCTION setCollectorVersionNumber(); + +CREATE TABLE collectorActiveVersions ( + id uuid primary key DEFAULT uuid_generate_v7(), + collectorId uuid not null, + versionId int not null, + foreign key (collectorId, versionId) references collectorVersions(collectorId, id) +); + CREATE TABLE collectorCodeVersions ( id uuid primary key DEFAULT uuid_generate_v7(), collectorId uuid not null, @@ -14,6 +44,8 @@ CREATE TABLE collectorCodeVersions ( minTextVersion int not null, addedVersion int not null, removedVersion int, + foreign key (collectorId, addedVersion) references collectorVersions(collectorId, id), + foreign key (collectorId, removedVersion) references collectorVersions(collectorId, id), foreign key (collectorId) references collectors(id) ); @@ -26,5 +58,7 @@ CREATE TABLE collectorQueries ( removedVersion int, foreign key (queryId) references queries(id), foreign key (collectorId) references collectors(id), + foreign key (collectorId, addedVersion) references collectorVersions(collectorId, id), + foreign key (collectorId, removedVersion) references collectorVersions(collectorId, id), unique (collectorId, name, removedVersion) ); \ No newline at end of file diff --git a/database/migrations/00000000000101_collector_views.up.sql b/database/migrations/00000000000101_collector_views.up.sql index c434ba11..f09e9310 100644 --- a/database/migrations/00000000000101_collector_views.up.sql +++ b/database/migrations/00000000000101_collector_views.up.sql @@ -1,23 +1,34 @@ +CREATE VIEW collectorCurrentActiveVersions as +SELECT DISTINCT + av.collectorId, + (FIRST_VALUE(av.versionId) OVER (PARTITION BY av.collectorId ORDER BY av.id DESC))::int as id + FROM collectorActiveVersions as av; + +CREATE VIEW collectorLatestVersions as + SELECT collectorId, max(id)::int as id + FROM collectorVersions + GROUP BY collectorId; + CREATE VIEW fullActiveCollectors AS - SELECT DISTINCT c.id, c.jobId, coalesce(cv.minCleanVersion, 0) as minCleanVersion, coalesce(cv.minTextVersion, 0) as minTextVersion, c.activeVersion, c.latestVersion, + SELECT DISTINCT c.id, c.jobId, coalesce(cv.minCleanVersion, 0) as minCleanVersion, coalesce(cv.minTextVersion, 0) as minTextVersion, coalesce(av.id, 0) as activeVersion, coalesce(lv.id, 0) as latestVersion, jsonb_object_agg(q.name, q.queryId) FILTER (WHERE q.name is not null) AS fields FROM collectors AS c + LEFT JOIN collectorCurrentActiveVersions as av on c.id = av.collectorId + LEFT JOIN collectorLatestVersions as lv on lv.collectorId = c.id LEFT JOIN collectorCodeVersions AS cv ON c.id = cv.collectorId - AND c.activeVersion >= cv.addedVersion - and c.activeVersion < COALESCE(cv.removedVersion, c.activeVersion + 1) + and isInVersion(av.id, cv.addedVersion, cv.removedVersion) LEFT JOIN collectorQueries AS q ON c.id = q.collectorId - AND c.activeVersion >= q.addedVersion - and c.activeVersion < COALESCE(q.removedVersion, c.activeVersion + 1) - GROUP BY c.id, c.jobId, cv.minCleanVersion, cv.minTextVersion; + and isInVersion(av.id, q.addedVersion, q.removedVersion) + GROUP BY c.id, c.jobId, cv.minCleanVersion, cv.minTextVersion, av.id, lv.id; CREATE VIEW collectorQueryDependencyTree AS WITH RECURSIVE coll as ( SELECT DISTINCT c.id, c.jobId, q.queryId FROM collectors as c + LEFT JOIN collectorCurrentActiveVersions as av on c.id = av.collectorId LEFT JOIN collectorQueries as q ON c.id = q.collectorId - AND c.activeVersion >= q.addedVersion - and c.activeVersion < COALESCE(q.removedVersion, c.activeVersion + 1) + and isInVersion(av.id, q.addedVersion, q.removedVersion) ), collectorQueryDependencyTree AS ( SELECT cq.id as collectorId, cq.jobId, aqc.id as queryId, aqc.type, aqc.requiredIds, aqc.activeVersion diff --git a/database/queries/collector.sql b/database/queries/collector.sql index 01a8e6cd..487c727e 100644 --- a/database/queries/collector.sql +++ b/database/queries/collector.sql @@ -10,8 +10,11 @@ SELECT * FROM fullActiveCollectors WHERE id = $1 LIMIT 1; -- name: CreateCollector :one INSERT INTO collectors (jobId) VALUES ($1) RETURNING id; --- name: UpdateCollector :exec -UPDATE collectors SET latestVersion = $1, activeVersion = $2 WHERE id = $3; +-- name: AddLatestCollectorVersion :exec +INSERT INTO collectorVersions (collectorId) VALUES ($1); + +-- name: AddActiveCollectorVersion :exec +INSERT INTO collectorActiveVersions (collectorId, versionId) VALUES ($1, $2); -- name: AddCollectorCodeVersion :exec INSERT INTO collectorCodeVersions (collectorId, addedVersion, minCleanVersion, minTextVersion) VALUES ($1, $2, $3, $4); diff --git a/database/queries/result.sql b/database/queries/result.sql index ea798d2c..a2615f8f 100644 --- a/database/queries/result.sql +++ b/database/queries/result.sql @@ -19,8 +19,9 @@ codeVersions as ( coalesce(ccv.minTextVersion, 1) as minTextVersion FROM docs as d LEFT JOIN collectors as c on c.jobId = d.jobId + LEFT JOIN collectorCurrentActiveVersions as av on c.id = av.collectorId LEFT JOIN collectorCodeVersions as ccv on c.id = ccv.collectorId - and isInVersion(c.activeVersion, ccv.addedVersion, ccv.removedVersion) + and isInVersion(av.id, ccv.addedVersion, ccv.removedVersion) LIMIT 1 ), latestVersions AS ( diff --git a/internal/database/repository/collector.sql.go b/internal/database/repository/collector.sql.go index 1c885b29..ae9e538a 100644 --- a/internal/database/repository/collector.sql.go +++ b/internal/database/repository/collector.sql.go @@ -11,6 +11,23 @@ import ( "github.com/jackc/pgx/v5/pgtype" ) +const addActiveCollectorVersion = `-- name: AddActiveCollectorVersion :exec +INSERT INTO collectorActiveVersions (collectorId, versionId) VALUES ($1, $2) +` + +type AddActiveCollectorVersionParams struct { + Collectorid pgtype.UUID `db:"collectorid"` + Versionid int32 `db:"versionid"` +} + +// AddActiveCollectorVersion +// +// INSERT INTO collectorActiveVersions (collectorId, versionId) VALUES ($1, $2) +func (q *Queries) AddActiveCollectorVersion(ctx context.Context, arg *AddActiveCollectorVersionParams) error { + _, err := q.db.Exec(ctx, addActiveCollectorVersion, arg.Collectorid, arg.Versionid) + return err +} + const addCollectorCodeVersion = `-- name: AddCollectorCodeVersion :exec INSERT INTO collectorCodeVersions (collectorId, addedVersion, minCleanVersion, minTextVersion) VALUES ($1, $2, $3, $4) ` @@ -59,6 +76,18 @@ func (q *Queries) AddCollectorQuery(ctx context.Context, arg *AddCollectorQueryP return err } +const addLatestCollectorVersion = `-- name: AddLatestCollectorVersion :exec +INSERT INTO collectorVersions (collectorId) VALUES ($1) +` + +// AddLatestCollectorVersion +// +// INSERT INTO collectorVersions (collectorId) VALUES ($1) +func (q *Queries) AddLatestCollectorVersion(ctx context.Context, collectorid pgtype.UUID) error { + _, err := q.db.Exec(ctx, addLatestCollectorVersion, collectorid) + return err +} + const createCollector = `-- name: CreateCollector :one INSERT INTO collectors (jobId) VALUES ($1) RETURNING id ` @@ -185,21 +214,3 @@ func (q *Queries) RemoveCollectorQuery(ctx context.Context, arg *RemoveCollector _, err := q.db.Exec(ctx, removeCollectorQuery, arg.Removedversion, arg.Queryid, arg.Collectorid) return err } - -const updateCollector = `-- name: UpdateCollector :exec -UPDATE collectors SET latestVersion = $1, activeVersion = $2 WHERE id = $3 -` - -type UpdateCollectorParams struct { - Latestversion int32 `db:"latestversion"` - Activeversion int32 `db:"activeversion"` - ID pgtype.UUID `db:"id"` -} - -// UpdateCollector -// -// UPDATE collectors SET latestVersion = $1, activeVersion = $2 WHERE id = $3 -func (q *Queries) UpdateCollector(ctx context.Context, arg *UpdateCollectorParams) error { - _, err := q.db.Exec(ctx, updateCollector, arg.Latestversion, arg.Activeversion, arg.ID) - return err -} diff --git a/internal/database/repository/collector_test.go b/internal/database/repository/collector_test.go index e259a762..e20ecc86 100644 --- a/internal/database/repository/collector_test.go +++ b/internal/database/repository/collector_test.go @@ -62,12 +62,43 @@ func TestCollector(t *testing.T) { Jobid: jobId, Mincleanversion: 0, Mintextversion: 0, - Activeversion: 1, - Latestversion: 1, + Activeversion: 0, + Latestversion: 0, }, coll) coll, err = queries.GetCollectorByJobID(ctx, jobId) assert.NoError(t, err) + assert.EqualExportedValues(t, &repository.Fullactivecollector{ + ID: collId, + Jobid: jobId, + Mincleanversion: 0, + Mintextversion: 0, + Activeversion: 0, + Latestversion: 0, + }, coll) + + err = queries.AddLatestCollectorVersion(ctx, collId) + assert.NoError(t, err) + + coll, err = queries.GetCollector(ctx, collId) + assert.NoError(t, err) + assert.EqualExportedValues(t, &repository.Fullactivecollector{ + ID: collId, + Jobid: jobId, + Mincleanversion: 0, + Mintextversion: 0, + Activeversion: 0, + Latestversion: 1, + }, coll) + + err = queries.AddActiveCollectorVersion(ctx, &repository.AddActiveCollectorVersionParams{ + Versionid: 1, + Collectorid: collId, + }) + assert.NoError(t, err) + + coll, err = queries.GetCollector(ctx, collId) + assert.NoError(t, err) assert.EqualExportedValues(t, &repository.Fullactivecollector{ ID: collId, Jobid: jobId, @@ -127,6 +158,9 @@ func TestCollector(t *testing.T) { }, }, qs) + err = queries.AddLatestCollectorVersion(ctx, collId) + assert.NoError(t, err) + removeV := int32(2) err = queries.RemoveCollectorQuery(ctx, &repository.RemoveCollectorQueryParams{ Collectorid: collId, @@ -141,10 +175,9 @@ func TestCollector(t *testing.T) { }) assert.NoError(t, err) - err = queries.UpdateCollector(ctx, &repository.UpdateCollectorParams{ - ID: collId, - Latestversion: 2, - Activeversion: 2, + err = queries.AddActiveCollectorVersion(ctx, &repository.AddActiveCollectorVersionParams{ + Versionid: 2, + Collectorid: collId, }) assert.NoError(t, err) diff --git a/internal/database/repository/models.go b/internal/database/repository/models.go index 0754928f..a6bca182 100644 --- a/internal/database/repository/models.go +++ b/internal/database/repository/models.go @@ -74,10 +74,14 @@ type Clientcansync struct { } type Collector struct { - ID pgtype.UUID `db:"id"` - Jobid pgtype.UUID `db:"jobid"` - Latestversion int32 `db:"latestversion"` - Activeversion int32 `db:"activeversion"` + ID pgtype.UUID `db:"id"` + Jobid pgtype.UUID `db:"jobid"` +} + +type Collectoractiveversion struct { + ID pgtype.UUID `db:"id"` + Collectorid pgtype.UUID `db:"collectorid"` + Versionid int32 `db:"versionid"` } type Collectorcodeversion struct { @@ -89,6 +93,16 @@ type Collectorcodeversion struct { Removedversion *int32 `db:"removedversion"` } +type Collectorcurrentactiveversion struct { + Collectorid pgtype.UUID `db:"collectorid"` + ID int32 `db:"id"` +} + +type Collectorlatestversion struct { + Collectorid pgtype.UUID `db:"collectorid"` + ID int32 `db:"id"` +} + type Collectorquery struct { ID pgtype.UUID `db:"id"` Collectorid pgtype.UUID `db:"collectorid"` @@ -107,6 +121,12 @@ type Collectorquerydependencytree struct { Requiredids []pgtype.UUID `db:"requiredids"` } +type Collectorversion struct { + Collectorid pgtype.UUID `db:"collectorid"` + ID int32 `db:"id"` + Addedat pgtype.Timestamp `db:"addedat"` +} + type Document struct { ID pgtype.UUID `db:"id"` Jobid pgtype.UUID `db:"jobid"` diff --git a/internal/database/repository/query_test.go b/internal/database/repository/query_test.go index 1151d90f..648eb6f9 100644 --- a/internal/database/repository/query_test.go +++ b/internal/database/repository/query_test.go @@ -246,6 +246,13 @@ func TestQueryDependencyTree(t *testing.T) { assert.NoError(t, err) collID, err := queries.CreateCollector(ctx, jobID) assert.NoError(t, err) + err = queries.AddLatestCollectorVersion(ctx, collID) + assert.NoError(t, err) + err = queries.AddActiveCollectorVersion(ctx, &repository.AddActiveCollectorVersionParams{ + Collectorid: collID, + Versionid: 1, + }) + assert.NoError(t, err) docID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Jobid: jobID, Hash: "sample", @@ -512,6 +519,13 @@ func TestListQueryJobs(t *testing.T) { assert.NoError(t, err) collOneID, err := queries.CreateCollector(ctx, jobOneID) assert.NoError(t, err) + err = queries.AddLatestCollectorVersion(ctx, collOneID) + assert.NoError(t, err) + err = queries.AddActiveCollectorVersion(ctx, &repository.AddActiveCollectorVersionParams{ + Collectorid: collOneID, + Versionid: 1, + }) + assert.NoError(t, err) err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Collectorid: collOneID, Queryid: contextID, @@ -528,6 +542,13 @@ func TestListQueryJobs(t *testing.T) { assert.NoError(t, err) collTwoID, err := queries.CreateCollector(ctx, jobTwoID) assert.NoError(t, err) + err = queries.AddLatestCollectorVersion(ctx, collTwoID) + assert.NoError(t, err) + err = queries.AddActiveCollectorVersion(ctx, &repository.AddActiveCollectorVersionParams{ + Collectorid: collTwoID, + Versionid: 1, + }) + assert.NoError(t, err) err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Collectorid: collTwoID, Queryid: contextID, @@ -546,6 +567,13 @@ func TestListQueryJobs(t *testing.T) { assert.NoError(t, err) collThreeID, err := queries.CreateCollector(ctx, jobThreeID) assert.NoError(t, err) + err = queries.AddLatestCollectorVersion(ctx, collThreeID) + assert.NoError(t, err) + err = queries.AddActiveCollectorVersion(ctx, &repository.AddActiveCollectorVersionParams{ + Collectorid: collThreeID, + Versionid: 1, + }) + assert.NoError(t, err) err = queries.AddCollectorQuery(ctx, &repository.AddCollectorQueryParams{ Collectorid: collThreeID, Queryid: contextID, diff --git a/internal/database/repository/result.sql.go b/internal/database/repository/result.sql.go index 9b4fa7de..7be7b9f2 100644 --- a/internal/database/repository/result.sql.go +++ b/internal/database/repository/result.sql.go @@ -65,8 +65,9 @@ codeVersions as ( coalesce(ccv.minTextVersion, 1) as minTextVersion FROM docs as d LEFT JOIN collectors as c on c.jobId = d.jobId + LEFT JOIN collectorCurrentActiveVersions as av on c.id = av.collectorId LEFT JOIN collectorCodeVersions as ccv on c.id = ccv.collectorId - and isInVersion(c.activeVersion, ccv.addedVersion, ccv.removedVersion) + and isInVersion(av.id, ccv.addedVersion, ccv.removedVersion) LIMIT 1 ), latestVersions AS ( @@ -132,8 +133,9 @@ type ListQueryRequirementValuesRow struct { // coalesce(ccv.minTextVersion, 1) as minTextVersion // FROM docs as d // LEFT JOIN collectors as c on c.jobId = d.jobId +// LEFT JOIN collectorCurrentActiveVersions as av on c.id = av.collectorId // LEFT JOIN collectorCodeVersions as ccv on c.id = ccv.collectorId -// and isInVersion(c.activeVersion, ccv.addedVersion, ccv.removedVersion) +// and isInVersion(av.id, ccv.addedVersion, ccv.removedVersion) // LIMIT 1 // ), // latestVersions AS ( diff --git a/internal/database/repository/result_test.go b/internal/database/repository/result_test.go index 37e797dd..134d97a6 100644 --- a/internal/database/repository/result_test.go +++ b/internal/database/repository/result_test.go @@ -292,6 +292,13 @@ func TestUnsyncedNoDepsQueries(t *testing.T) { assert.NoError(t, err) collectorId, err := queries.CreateCollector(ctx, jobId) assert.NoError(t, err) + err = queries.AddLatestCollectorVersion(ctx, collectorId) + assert.NoError(t, err) + err = queries.AddActiveCollectorVersion(ctx, &repository.AddActiveCollectorVersionParams{ + Collectorid: collectorId, + Versionid: 1, + }) + assert.NoError(t, err) documentID, err := queries.CreateDocument(ctx, &repository.CreateDocumentParams{ Jobid: jobId, Hash: "example_hash", diff --git a/internal/job/collector/create.go b/internal/job/collector/create.go index 8fe4438c..4113680a 100644 --- a/internal/job/collector/create.go +++ b/internal/job/collector/create.go @@ -110,14 +110,27 @@ func (s *Service) submitCreate(ctx context.Context, params *dbCreateParams) (uui return err } dbID = dID - latestVersion := int32(1) + + err = qtx.AddLatestCollectorVersion(ctx, dbID) + if err != nil { + return err + } + version := int32(1) + + err = qtx.AddActiveCollectorVersion(ctx, &repository.AddActiveCollectorVersionParams{ + Collectorid: dbID, + Versionid: version, + }) + if err != nil { + return err + } if params.MinCleanVersion != nil || params.MinTextVersion != nil { err = qtx.AddCollectorCodeVersion(ctx, &repository.AddCollectorCodeVersionParams{ Collectorid: dbID, Mincleanversion: validation.GetUpdatedValue(1, params.MinCleanVersion), Mintextversion: validation.GetUpdatedValue(1, params.MinTextVersion), - Addedversion: latestVersion, + Addedversion: version, }) if err != nil { return err @@ -130,7 +143,7 @@ func (s *Service) submitCreate(ctx context.Context, params *dbCreateParams) (uui Collectorid: dbID, Name: key, Queryid: field, - Addedversion: latestVersion, + Addedversion: version, }) if err != nil { return err diff --git a/internal/job/collector/create_test.go b/internal/job/collector/create_test.go index be4aad76..646928c3 100644 --- a/internal/job/collector/create_test.go +++ b/internal/job/collector/create_test.go @@ -48,6 +48,10 @@ func TestCreate(t *testing.T) { pgxmock.NewRows([]string{"id"}). AddRow(database.MustToDBUUID(id)), ) + pool.ExpectExec("name: AddLatestCollectorVersion :exec").WithArgs(database.MustToDBUUID(id)). + WillReturnResult(pgxmock.NewResult("", 1)) + pool.ExpectExec("name: AddActiveCollectorVersion :exec").WithArgs(database.MustToDBUUID(id), int32(1)). + WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: AddCollectorCodeVersion :exec").WithArgs(database.MustToDBUUID(id), int32(1), *create.MinCleanVersion, *create.MinTextVersion). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(id), "example_key", database.MustToDBUUID((*create.Fields)["example_key"]), int32(1)). diff --git a/internal/job/collector/createprivate_test.go b/internal/job/collector/createprivate_test.go index d3bc6055..0231b2b7 100644 --- a/internal/job/collector/createprivate_test.go +++ b/internal/job/collector/createprivate_test.go @@ -96,6 +96,10 @@ func TestSubmitCreate(t *testing.T) { pgxmock.NewRows([]string{"id"}). AddRow(database.MustToDBUUID(id)), ) + pool.ExpectExec("name: AddLatestCollectorVersion :exec").WithArgs(database.MustToDBUUID(id)). + WillReturnResult(pgxmock.NewResult("", 1)) + pool.ExpectExec("name: AddActiveCollectorVersion :exec").WithArgs(database.MustToDBUUID(id), int32(1)). + WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: AddCollectorCodeVersion :exec").WithArgs(database.MustToDBUUID(id), int32(1), *params.MinCleanVersion, *params.MinTextVersion). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(id), "example_key", (*params.Fields)["example_key"], int32(1)). diff --git a/internal/job/collector/update/update.go b/internal/job/collector/update/update.go index 5f16daf4..5ca0f8eb 100644 --- a/internal/job/collector/update/update.go +++ b/internal/job/collector/update/update.go @@ -186,9 +186,33 @@ func (s *Service) normalizeUpdateFieldsToDB(ctx context.Context, current map[str func (s *Service) submitUpdate(ctx context.Context, current *collector.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) + activeName, err := validation.GetFieldName(params, params.ActiveVersion) + if err != nil { + return err + } + onlyactive := validation.AreAllPointersNilExcept(params, activeName) + + if !onlyactive { + err := qtx.AddLatestCollectorVersion(ctx, id) + if err != nil { + return err + } + } + + if params.ActiveVersion != nil { + err := qtx.AddActiveCollectorVersion(ctx, &repository.AddActiveCollectorVersionParams{ + Collectorid: id, + Versionid: *params.ActiveVersion, + }) + if err != nil { + return err + } + } + + latestVersion := current.LatestVersion + 1 + if params.MinCleanVersion != nil || params.MinTextVersion != nil { err := qtx.RemoveCollectorCodeVersion(ctx, &repository.RemoveCollectorCodeVersionParams{ Collectorid: id, @@ -236,30 +260,6 @@ func (s *Service) submitUpdate(ctx context.Context, current *collector.Collector } } - activeVersion := params.ActiveVersion - if activeVersion == nil { - activeVersion = ¤t.ActiveVersion - } - - activeName, err := validation.GetFieldName(params, params.ActiveVersion) - if err != nil { - return err - } - onlyactive := validation.AreAllPointersNilExcept(params, activeName) - - if onlyactive { - latestVersion-- - } - - err = qtx.UpdateCollector(ctx, &repository.UpdateCollectorParams{ - ID: id, - Latestversion: latestVersion, - Activeversion: *activeVersion, - }) - if err != nil { - return err - } - slog.Debug("job collector updated", "update", *params) return nil diff --git a/internal/job/collector/update/update_test.go b/internal/job/collector/update/update_test.go index 068b55a7..e32553f9 100644 --- a/internal/job/collector/update/update_test.go +++ b/internal/job/collector/update/update_test.go @@ -64,12 +64,14 @@ func TestUpdate(t *testing.T) { ) pool.ExpectBeginTx(pgx.TxOptions{}) rv := current.LatestVersion + 1 + pool.ExpectExec("name: AddLatestCollectorVersion :exec").WithArgs(database.MustToDBUUID(current.ID)). + WillReturnResult(pgxmock.NewResult("", 1)) + pool.ExpectExec("name: AddActiveCollectorVersion :exec").WithArgs(database.MustToDBUUID(current.ID), int32(2)). + WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: RemoveCollectorCodeVersion :exec").WithArgs(&rv, database.MustToDBUUID(current.ID)). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: AddCollectorCodeVersion :exec").WithArgs(database.MustToDBUUID(current.ID), int32(5), *update.MinCleanVersion, int32(0)). WillReturnResult(pgxmock.NewResult("", 1)) - pool.ExpectExec("name: UpdateCollector :exec").WithArgs(int32(5), int32(2), database.MustToDBUUID(current.ID)). - WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() mockSQS.EXPECT(). diff --git a/internal/job/collector/update/updateprivate_test.go b/internal/job/collector/update/updateprivate_test.go index 5364628a..e0c8854d 100644 --- a/internal/job/collector/update/updateprivate_test.go +++ b/internal/job/collector/update/updateprivate_test.go @@ -219,6 +219,10 @@ func TestSubmitUpdate(t *testing.T) { pool.ExpectBeginTx(pgx.TxOptions{}) rv := int32(2) + pool.ExpectExec("name: AddLatestCollectorVersion :exec").WithArgs(database.MustToDBUUID(current.ID)). + WillReturnResult(pgxmock.NewResult("", 1)) + pool.ExpectExec("name: AddActiveCollectorVersion :exec").WithArgs(database.MustToDBUUID(current.ID), int32(2)). + WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: RemoveCollectorCodeVersion :exec").WithArgs(&rv, database.MustToDBUUID(current.ID)). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: AddCollectorCodeVersion :exec").WithArgs(database.MustToDBUUID(current.ID), int32(2), *params.MinCleanVersion, *params.MinTextVersion). @@ -232,8 +236,6 @@ func TestSubmitUpdate(t *testing.T) { WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectExec("name: AddCollectorQuery :exec").WithArgs(database.MustToDBUUID(current.ID), "changed_key", (*params.Fields)["changed_key"], int32(2)). WillReturnResult(pgxmock.NewResult("", 1)) - pool.ExpectExec("name: UpdateCollector :exec").WithArgs(int32(2), int32(2), database.MustToDBUUID(current.ID)). - WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() err = svc.submitUpdate(ctx, ¤t, ¶ms) @@ -272,7 +274,7 @@ func TestSubmitUpdate(t *testing.T) { } pool.ExpectBeginTx(pgx.TxOptions{}) - pool.ExpectExec("name: UpdateCollector :exec").WithArgs(int32(1), av, database.MustToDBUUID(current.ID)). + pool.ExpectExec("name: AddActiveCollectorVersion :exec").WithArgs(database.MustToDBUUID(current.ID), int32(2)). WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() diff --git a/internal/job/create_test.go b/internal/job/create_test.go index a18ee0cc..d01b8689 100644 --- a/internal/job/create_test.go +++ b/internal/job/create_test.go @@ -34,6 +34,7 @@ func TestCreate(t *testing.T) { ID: uuid.New(), ClientID: uuid.New(), } + collId := uuid.New() pool.ExpectQuery("name: CreateJob :one").WithArgs(database.MustToDBUUID(job.ClientID)).WillReturnRows( pgxmock.NewRows([]string{"id"}). @@ -42,8 +43,12 @@ func TestCreate(t *testing.T) { pool.ExpectBeginTx(pgx.TxOptions{}) pool.ExpectQuery("name: CreateCollector :one").WithArgs(database.MustToDBUUID(job.ID)).WillReturnRows( pgxmock.NewRows([]string{"id"}). - AddRow(database.MustToDBUUID(uuid.New())), + AddRow(database.MustToDBUUID(collId)), ) + pool.ExpectExec("name: AddLatestCollectorVersion :exec").WithArgs(database.MustToDBUUID(collId)). + WillReturnResult(pgxmock.NewResult("", 1)) + pool.ExpectExec("name: AddActiveCollectorVersion :exec").WithArgs(database.MustToDBUUID(collId), int32(1)). + WillReturnResult(pgxmock.NewResult("", 1)) pool.ExpectCommit() aid, err := svc.Create(ctx, job.ClientID)