@@ -66,6 +66,10 @@ func TestCreateQuery(t *testing.T) {
|
||||
pgxmock.NewRows([]string{"id"}).
|
||||
AddRow(database.MustToDBUUID(id)),
|
||||
)
|
||||
pool.ExpectExec("name: AddLatestQueryVersion :exec").WithArgs(database.MustToDBUUID(id)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(id), int32(1)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectCommit()
|
||||
|
||||
err = cons.CreateQuery(ctx)
|
||||
@@ -203,7 +207,7 @@ func TestUpdateQuery(t *testing.T) {
|
||||
)
|
||||
|
||||
pool.ExpectBeginTx(pgx.TxOptions{})
|
||||
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(av, int32(2), database.MustToDBUUID(id)).
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(id), av).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectCommit()
|
||||
mockSQS.EXPECT().
|
||||
|
||||
@@ -17,4 +17,18 @@ select encode(
|
||||
'hex')::uuid;
|
||||
$$
|
||||
language SQL
|
||||
volatile;
|
||||
volatile;
|
||||
|
||||
CREATE FUNCTION isInVersion(
|
||||
_version int,
|
||||
_addedVersion int,
|
||||
_removedVersion int
|
||||
)
|
||||
RETURNS boolean as $$
|
||||
BEGIN
|
||||
RETURN _version is null OR (
|
||||
_version >= _addedVersion AND
|
||||
(_removedVersion is null OR _version < _removedVersion)
|
||||
);
|
||||
END;
|
||||
$$ LANGUAGE plpgsql;
|
||||
@@ -2,11 +2,41 @@ CREATE TYPE queryType AS ENUM ('context_full', 'json_extractor');
|
||||
|
||||
CREATE TABLE queries (
|
||||
id uuid primary key DEFAULT uuid_generate_v7(),
|
||||
latestVersion int not null default 1,
|
||||
activeVersion int not null default 1,
|
||||
type queryType not null
|
||||
);
|
||||
|
||||
CREATE TABLE queryVersions (
|
||||
queryId uuid not null,
|
||||
id int not null,
|
||||
addedAt timestamp not null default current_timestamp,
|
||||
primary key (id, queryId),
|
||||
foreign key (queryId) references queries(id)
|
||||
);
|
||||
|
||||
CREATE OR REPLACE FUNCTION setQueryVersionNumber()
|
||||
RETURNS TRIGGER AS $$
|
||||
BEGIN
|
||||
SELECT COALESCE(MAX(id), 0) + 1
|
||||
INTO NEW.id
|
||||
FROM queryVersions
|
||||
WHERE queryId = NEW.queryId;
|
||||
|
||||
RETURN NEW;
|
||||
END;
|
||||
$$ LANGUAGE plpgsql;
|
||||
|
||||
CREATE TRIGGER setVersionNumberTrigger
|
||||
BEFORE INSERT ON queryVersions
|
||||
FOR EACH ROW
|
||||
EXECUTE FUNCTION setQueryVersionNumber();
|
||||
|
||||
CREATE TABLE queryActiveVersions (
|
||||
id uuid primary key DEFAULT uuid_generate_v7(),
|
||||
queryId uuid not null,
|
||||
versionId int not null,
|
||||
foreign key (queryId, versionId) references queryVersions(queryId, id)
|
||||
);
|
||||
|
||||
CREATE TABLE requiredQueries (
|
||||
id uuid primary key DEFAULT uuid_generate_v7(),
|
||||
queryId uuid not null,
|
||||
@@ -15,6 +45,8 @@ CREATE TABLE requiredQueries (
|
||||
removedVersion int,
|
||||
foreign key (queryId) references queries(id),
|
||||
foreign key (requiredQueryId) references queries(id),
|
||||
foreign key (queryId, addedVersion) references queryVersions(queryId, id),
|
||||
foreign key (queryId, removedVersion) references queryVersions(queryId, id),
|
||||
unique (queryId, requiredQueryId, removedVersion)
|
||||
);
|
||||
|
||||
@@ -25,5 +57,7 @@ CREATE TABLE queryConfigs (
|
||||
addedVersion int not null,
|
||||
removedVersion int,
|
||||
foreign key (queryId) references queries(id),
|
||||
foreign key (queryId, addedVersion) references queryVersions(queryId, id),
|
||||
foreign key (queryId, removedVersion) references queryVersions(queryId, id),
|
||||
unique (queryId, removedVersion)
|
||||
);
|
||||
@@ -8,6 +8,5 @@ CREATE TABLE clientCanSync (
|
||||
id uuid primary key DEFAULT uuid_generate_v7(),
|
||||
clientId uuid not null,
|
||||
canSync boolean not null,
|
||||
addedAt timestamp not null default current_timestamp,
|
||||
foreign key (clientId) references clients(id)
|
||||
);
|
||||
@@ -8,6 +8,5 @@ CREATE TABLE jobCanSync (
|
||||
id uuid primary key DEFAULT uuid_generate_v7(),
|
||||
jobId uuid not null,
|
||||
canSync boolean not null,
|
||||
addedAt timestamp not null default current_timestamp,
|
||||
foreign key (jobId) references jobs(id)
|
||||
);
|
||||
@@ -8,5 +8,6 @@ CREATE TABLE results (
|
||||
queryVersion int not null,
|
||||
foreign key (queryId) references queries(id),
|
||||
foreign key (documentId) references documents(id),
|
||||
unique (queryId, documentId, cleanVersion, textVersion, queryVersion)
|
||||
unique (queryId, documentId, cleanVersion, textVersion, queryVersion),
|
||||
foreign key (queryId, queryVersion) references queryVersions(queryId, id)
|
||||
);
|
||||
@@ -1,10 +1,21 @@
|
||||
CREATE VIEW queryCurrentActiveVersions as
|
||||
SELECT DISTINCT
|
||||
av.queryId,
|
||||
(FIRST_VALUE(av.versionId) OVER (PARTITION BY av.queryId ORDER BY av.id DESC))::int as id
|
||||
FROM queryActiveVersions as av;
|
||||
|
||||
CREATE VIEW queryLatestVersions as
|
||||
SELECT queryId, max(id)::int as id
|
||||
FROM queryVersions
|
||||
GROUP BY queryId;
|
||||
|
||||
CREATE VIEW fullActiveQueries AS
|
||||
WITH config as (
|
||||
with config as (
|
||||
SELECT c.queryId, c.config
|
||||
FROM queries AS q
|
||||
LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
LEFT JOIN queryConfigs AS c ON q.id = c.queryId
|
||||
and q.activeVersion >= c.addedVersion
|
||||
and (c.removedVersion is null or q.activeVersion < c.removedVersion)
|
||||
and isInVersion(av.id, c.addedVersion, c.removedVersion)
|
||||
),
|
||||
requiredIds as (
|
||||
SELECT r.queryId,
|
||||
@@ -12,35 +23,37 @@ requiredIds as (
|
||||
FILTER (WHERE r.requiredQueryId != '00000000-0000-0000-0000-000000000000')::uuid[]
|
||||
as requiredIds
|
||||
FROM queries AS q
|
||||
LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
LEFT JOIN requiredQueries AS r ON q.id = r.queryId
|
||||
and q.activeVersion >= r.addedVersion
|
||||
and (r.removedVersion is null or q.activeVersion < r.removedVersion)
|
||||
and isInVersion(av.id, r.addedVersion, r.removedVersion)
|
||||
GROUP BY r.queryId
|
||||
)
|
||||
SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, c.config,
|
||||
SELECT DISTINCT q.id, q.type, coalesce(av.id, 0) as activeVersion, coalesce(lv.id, 0) as latestVersion, c.config,
|
||||
coalesce(
|
||||
r.requiredIds,
|
||||
array[]::uuid[]
|
||||
)::uuid[] as requiredIds
|
||||
FROM queries AS q
|
||||
LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
LEFT JOIN queryLatestVersions as lv on lv.queryId = q.id
|
||||
LEFT JOIN config AS c ON q.id = c.queryId
|
||||
LEFT JOIN requiredIds AS r ON q.id = r.queryId;
|
||||
|
||||
CREATE VIEW queryActiveDependencies AS
|
||||
WITH RECURSIVE queryActiveDependencies AS (
|
||||
with recursive queryActiveDependencies AS (
|
||||
SELECT q.id, r.requiredQueryId
|
||||
FROM queries AS q
|
||||
LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
LEFT JOIN requiredQueries AS r ON q.id = r.queryId
|
||||
and q.activeVersion >= r.addedVersion
|
||||
and q.activeVersion < COALESCE(r.removedVersion, q.activeVersion + 1)
|
||||
and isInVersion(av.id, r.addedVersion, r.removedVersion)
|
||||
|
||||
UNION ALL
|
||||
|
||||
SELECT qd.id, r.requiredQueryId
|
||||
FROM queries AS q
|
||||
LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
LEFT JOIN requiredQueries AS r ON q.id = r.queryId
|
||||
and q.activeVersion >= r.addedVersion
|
||||
and q.activeVersion < COALESCE(r.removedVersion, q.activeVersion + 1)
|
||||
and isInVersion(av.id, r.addedVersion, r.removedVersion)
|
||||
JOIN queryActiveDependencies as qd on q.id = qd.requiredQueryId
|
||||
)
|
||||
SELECT DISTINCT id, requiredQueryId FROM queryActiveDependencies;
|
||||
@@ -8,7 +8,7 @@ SELECT c.id, c.name, coalesce(cs.canSync, false) as canSync
|
||||
SELECT clientId, canSync
|
||||
FROM clientCanSync
|
||||
WHERE clientId = $1
|
||||
ORDER BY addedAt DESC
|
||||
ORDER BY id DESC
|
||||
LIMIT 1
|
||||
) as cs on cs.clientId = c.id
|
||||
WHERE c.id = $1;
|
||||
|
||||
@@ -5,7 +5,7 @@ SELECT j.id, j.clientId, coalesce(cs.canSync, false) as canSync
|
||||
SELECT jobId, canSync
|
||||
FROM jobCanSync
|
||||
WHERE jobId = $1
|
||||
ORDER BY addedAt DESC
|
||||
ORDER BY id DESC
|
||||
LIMIT 1
|
||||
) as cs on cs.jobId = j.id
|
||||
WHERE j.id = $1;
|
||||
|
||||
@@ -6,14 +6,13 @@ SELECT * FROM fullActiveQueries WHERE id = $1;
|
||||
|
||||
-- name: GetQueryWithVersion :one
|
||||
WITH query as (
|
||||
SELECT id, type, activeVersion, latestVersion FROM queries WHERE id = @id
|
||||
SELECT id, type FROM queries WHERE id = @id
|
||||
),
|
||||
config as (
|
||||
SELECT c.queryId, c.config
|
||||
FROM query AS q
|
||||
LEFT JOIN queryConfigs AS c ON q.id = c.queryId
|
||||
and @version >= c.addedVersion
|
||||
and (c.removedVersion is null or @version < c.removedVersion)
|
||||
and isInVersion(@version, c.addedVersion, c.removedVersion)
|
||||
),
|
||||
requiredIds as (
|
||||
SELECT r.queryId,
|
||||
@@ -22,16 +21,17 @@ requiredIds as (
|
||||
as requiredIds
|
||||
FROM query AS q
|
||||
LEFT JOIN requiredQueries AS r ON q.id = r.queryId
|
||||
and @version >= r.addedVersion
|
||||
and (r.removedVersion is null or @version < r.removedVersion)
|
||||
and isInVersion(@version, r.addedVersion, r.removedVersion)
|
||||
GROUP BY r.queryId
|
||||
)
|
||||
SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, c.config,
|
||||
SELECT DISTINCT q.id, q.type, coalesce(av.id, 0) as activeVersion, coalesce(lv.id, 0) as latestVersion, c.config,
|
||||
coalesce(
|
||||
r.requiredIds,
|
||||
array[]::uuid[]
|
||||
)::uuid[] as requiredIds
|
||||
FROM query AS q
|
||||
LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
LEFT JOIN queryLatestVersions as lv on lv.queryId = q.id
|
||||
LEFT JOIN config AS c ON q.id = c.queryId
|
||||
LEFT JOIN requiredIds AS r ON q.id = r.queryId;
|
||||
|
||||
@@ -44,8 +44,11 @@ SELECT * FROM fullActiveQueries WHERE id = any($1);
|
||||
-- name: CreateQuery :one
|
||||
INSERT INTO queries (type) VALUES ($1) RETURNING id;
|
||||
|
||||
-- name: UpdateQuery :exec
|
||||
UPDATE queries SET activeVersion = $1, latestVersion = $2 WHERE id = $3;
|
||||
-- name: AddLatestQueryVersion :exec
|
||||
INSERT INTO queryVersions (queryId) VALUES ($1);
|
||||
|
||||
-- name: AddActiveQueryVersion :exec
|
||||
INSERT INTO queryActiveVersions (queryId, versionId) VALUES ($1, $2);
|
||||
|
||||
-- name: AddRequiredQuery :exec
|
||||
INSERT INTO requiredQueries (queryId, requiredQueryId, addedVersion) VALUES ($1, $2, $3);
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
-- name: ListQueryRequirementValues :many
|
||||
WITH reqQueries as (
|
||||
SELECT q.id as queryId, q.activeVersion, q.type
|
||||
SELECT q.id as queryId, av.id as activeVersion, q.type
|
||||
FROM requiredQueries as rq
|
||||
JOIN queries as q on q.id = rq.requiredQueryId
|
||||
LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
WHERE rq.queryId = @queryId
|
||||
and @version >= rq.addedVersion
|
||||
and (rq.removedVersion is null or @version < rq.removedVersion)
|
||||
and isInVersion(@version, rq.addedVersion, rq.removedVersion)
|
||||
),
|
||||
docs as (
|
||||
SELECT id, jobId
|
||||
@@ -20,8 +20,7 @@ codeVersions as (
|
||||
FROM docs as d
|
||||
LEFT JOIN collectors as c on c.jobId = d.jobId
|
||||
LEFT JOIN collectorCodeVersions as ccv on c.id = ccv.collectorId
|
||||
and c.activeVersion >= ccv.addedVersion
|
||||
and c.activeVersion < COALESCE(ccv.removedVersion, c.activeVersion)
|
||||
and isInVersion(c.activeVersion, ccv.addedVersion, ccv.removedVersion)
|
||||
LIMIT 1
|
||||
),
|
||||
latestVersions AS (
|
||||
|
||||
@@ -49,7 +49,7 @@ SELECT c.id, c.name, coalesce(cs.canSync, false) as canSync
|
||||
SELECT clientId, canSync
|
||||
FROM clientCanSync
|
||||
WHERE clientId = $1
|
||||
ORDER BY addedAt DESC
|
||||
ORDER BY id DESC
|
||||
LIMIT 1
|
||||
) as cs on cs.clientId = c.id
|
||||
WHERE c.id = $1
|
||||
@@ -69,7 +69,7 @@ type GetClientRow struct {
|
||||
// SELECT clientId, canSync
|
||||
// FROM clientCanSync
|
||||
// WHERE clientId = $1
|
||||
// ORDER BY addedAt DESC
|
||||
// ORDER BY id DESC
|
||||
// LIMIT 1
|
||||
// ) as cs on cs.clientId = c.id
|
||||
// WHERE c.id = $1
|
||||
|
||||
@@ -36,6 +36,8 @@ func TestCollector(t *testing.T) {
|
||||
assert.NoError(t, err)
|
||||
jsonId, err := queries.CreateQuery(ctx, repository.QuerytypeJsonExtractor)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddLatestQueryVersion(ctx, jsonId)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
|
||||
Queryid: jsonId,
|
||||
Requiredqueryid: contextId,
|
||||
@@ -111,7 +113,7 @@ func TestCollector(t *testing.T) {
|
||||
Collectorid: collId,
|
||||
Jobid: jobId,
|
||||
Queryid: jsonId,
|
||||
Queryversion: 1,
|
||||
Queryversion: 0,
|
||||
Type: repository.QuerytypeJsonExtractor,
|
||||
Requiredids: []pgtype.UUID{contextId},
|
||||
},
|
||||
@@ -119,7 +121,7 @@ func TestCollector(t *testing.T) {
|
||||
Collectorid: collId,
|
||||
Jobid: jobId,
|
||||
Queryid: contextId,
|
||||
Queryversion: 1,
|
||||
Queryversion: 0,
|
||||
Type: repository.QuerytypeContextFull,
|
||||
Requiredids: []pgtype.UUID{},
|
||||
},
|
||||
|
||||
@@ -49,7 +49,7 @@ SELECT j.id, j.clientId, coalesce(cs.canSync, false) as canSync
|
||||
SELECT jobId, canSync
|
||||
FROM jobCanSync
|
||||
WHERE jobId = $1
|
||||
ORDER BY addedAt DESC
|
||||
ORDER BY id DESC
|
||||
LIMIT 1
|
||||
) as cs on cs.jobId = j.id
|
||||
WHERE j.id = $1
|
||||
@@ -69,7 +69,7 @@ type GetJobRow struct {
|
||||
// SELECT jobId, canSync
|
||||
// FROM jobCanSync
|
||||
// WHERE jobId = $1
|
||||
// ORDER BY addedAt DESC
|
||||
// ORDER BY id DESC
|
||||
// LIMIT 1
|
||||
// ) as cs on cs.jobId = j.id
|
||||
// WHERE j.id = $1
|
||||
|
||||
@@ -68,10 +68,9 @@ type Client struct {
|
||||
}
|
||||
|
||||
type Clientcansync struct {
|
||||
ID pgtype.UUID `db:"id"`
|
||||
Clientid pgtype.UUID `db:"clientid"`
|
||||
Cansync bool `db:"cansync"`
|
||||
Addedat pgtype.Timestamp `db:"addedat"`
|
||||
ID pgtype.UUID `db:"id"`
|
||||
Clientid pgtype.UUID `db:"clientid"`
|
||||
Cansync bool `db:"cansync"`
|
||||
}
|
||||
|
||||
type Collector struct {
|
||||
@@ -162,17 +161,14 @@ type Job struct {
|
||||
}
|
||||
|
||||
type Jobcansync struct {
|
||||
ID pgtype.UUID `db:"id"`
|
||||
Jobid pgtype.UUID `db:"jobid"`
|
||||
Cansync bool `db:"cansync"`
|
||||
Addedat pgtype.Timestamp `db:"addedat"`
|
||||
ID pgtype.UUID `db:"id"`
|
||||
Jobid pgtype.UUID `db:"jobid"`
|
||||
Cansync bool `db:"cansync"`
|
||||
}
|
||||
|
||||
type Query struct {
|
||||
ID pgtype.UUID `db:"id"`
|
||||
Latestversion int32 `db:"latestversion"`
|
||||
Activeversion int32 `db:"activeversion"`
|
||||
Type Querytype `db:"type"`
|
||||
ID pgtype.UUID `db:"id"`
|
||||
Type Querytype `db:"type"`
|
||||
}
|
||||
|
||||
type Queryactivedependency struct {
|
||||
@@ -180,6 +176,12 @@ type Queryactivedependency struct {
|
||||
Requiredqueryid pgtype.UUID `db:"requiredqueryid"`
|
||||
}
|
||||
|
||||
type Queryactiveversion struct {
|
||||
ID pgtype.UUID `db:"id"`
|
||||
Queryid pgtype.UUID `db:"queryid"`
|
||||
Versionid int32 `db:"versionid"`
|
||||
}
|
||||
|
||||
type Queryconfig struct {
|
||||
ID pgtype.UUID `db:"id"`
|
||||
Queryid pgtype.UUID `db:"queryid"`
|
||||
@@ -188,6 +190,22 @@ type Queryconfig struct {
|
||||
Removedversion *int32 `db:"removedversion"`
|
||||
}
|
||||
|
||||
type Querycurrentactiveversion struct {
|
||||
Queryid pgtype.UUID `db:"queryid"`
|
||||
ID int32 `db:"id"`
|
||||
}
|
||||
|
||||
type Querylatestversion struct {
|
||||
Queryid pgtype.UUID `db:"queryid"`
|
||||
ID int32 `db:"id"`
|
||||
}
|
||||
|
||||
type Queryversion struct {
|
||||
Queryid pgtype.UUID `db:"queryid"`
|
||||
ID int32 `db:"id"`
|
||||
Addedat pgtype.Timestamp `db:"addedat"`
|
||||
}
|
||||
|
||||
type Requiredquery struct {
|
||||
ID pgtype.UUID `db:"id"`
|
||||
Queryid pgtype.UUID `db:"queryid"`
|
||||
|
||||
@@ -11,6 +11,35 @@ import (
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
)
|
||||
|
||||
const addActiveQueryVersion = `-- name: AddActiveQueryVersion :exec
|
||||
INSERT INTO queryActiveVersions (queryId, versionId) VALUES ($1, $2)
|
||||
`
|
||||
|
||||
type AddActiveQueryVersionParams struct {
|
||||
Queryid pgtype.UUID `db:"queryid"`
|
||||
Versionid int32 `db:"versionid"`
|
||||
}
|
||||
|
||||
// AddActiveQueryVersion
|
||||
//
|
||||
// INSERT INTO queryActiveVersions (queryId, versionId) VALUES ($1, $2)
|
||||
func (q *Queries) AddActiveQueryVersion(ctx context.Context, arg *AddActiveQueryVersionParams) error {
|
||||
_, err := q.db.Exec(ctx, addActiveQueryVersion, arg.Queryid, arg.Versionid)
|
||||
return err
|
||||
}
|
||||
|
||||
const addLatestQueryVersion = `-- name: AddLatestQueryVersion :exec
|
||||
INSERT INTO queryVersions (queryId) VALUES ($1)
|
||||
`
|
||||
|
||||
// AddLatestQueryVersion
|
||||
//
|
||||
// INSERT INTO queryVersions (queryId) VALUES ($1)
|
||||
func (q *Queries) AddLatestQueryVersion(ctx context.Context, queryid pgtype.UUID) error {
|
||||
_, err := q.db.Exec(ctx, addLatestQueryVersion, queryid)
|
||||
return err
|
||||
}
|
||||
|
||||
const addQueryConfig = `-- name: AddQueryConfig :exec
|
||||
INSERT INTO queryConfigs (queryId, config, addedVersion) VALUES ($1, $2, $3)
|
||||
`
|
||||
@@ -126,14 +155,13 @@ func (q *Queries) GetQueryConfig(ctx context.Context, arg *GetQueryConfigParams)
|
||||
|
||||
const getQueryWithVersion = `-- name: GetQueryWithVersion :one
|
||||
WITH query as (
|
||||
SELECT id, type, activeVersion, latestVersion FROM queries WHERE id = $1
|
||||
SELECT id, type FROM queries WHERE id = $1
|
||||
),
|
||||
config as (
|
||||
SELECT c.queryId, c.config
|
||||
FROM query AS q
|
||||
LEFT JOIN queryConfigs AS c ON q.id = c.queryId
|
||||
and $2 >= c.addedVersion
|
||||
and (c.removedVersion is null or $2 < c.removedVersion)
|
||||
and isInVersion($2, c.addedVersion, c.removedVersion)
|
||||
),
|
||||
requiredIds as (
|
||||
SELECT r.queryId,
|
||||
@@ -142,16 +170,17 @@ requiredIds as (
|
||||
as requiredIds
|
||||
FROM query AS q
|
||||
LEFT JOIN requiredQueries AS r ON q.id = r.queryId
|
||||
and $2 >= r.addedVersion
|
||||
and (r.removedVersion is null or $2 < r.removedVersion)
|
||||
and isInVersion($2, r.addedVersion, r.removedVersion)
|
||||
GROUP BY r.queryId
|
||||
)
|
||||
SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, c.config,
|
||||
SELECT DISTINCT q.id, q.type, coalesce(av.id, 0) as activeVersion, coalesce(lv.id, 0) as latestVersion, c.config,
|
||||
coalesce(
|
||||
r.requiredIds,
|
||||
array[]::uuid[]
|
||||
)::uuid[] as requiredIds
|
||||
FROM query AS q
|
||||
LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
LEFT JOIN queryLatestVersions as lv on lv.queryId = q.id
|
||||
LEFT JOIN config AS c ON q.id = c.queryId
|
||||
LEFT JOIN requiredIds AS r ON q.id = r.queryId
|
||||
`
|
||||
@@ -173,14 +202,13 @@ type GetQueryWithVersionRow struct {
|
||||
// GetQueryWithVersion
|
||||
//
|
||||
// WITH query as (
|
||||
// SELECT id, type, activeVersion, latestVersion FROM queries WHERE id = $1
|
||||
// SELECT id, type FROM queries WHERE id = $1
|
||||
// ),
|
||||
// config as (
|
||||
// SELECT c.queryId, c.config
|
||||
// FROM query AS q
|
||||
// LEFT JOIN queryConfigs AS c ON q.id = c.queryId
|
||||
// and $2 >= c.addedVersion
|
||||
// and (c.removedVersion is null or $2 < c.removedVersion)
|
||||
// and isInVersion($2, c.addedVersion, c.removedVersion)
|
||||
// ),
|
||||
// requiredIds as (
|
||||
// SELECT r.queryId,
|
||||
@@ -189,16 +217,17 @@ type GetQueryWithVersionRow struct {
|
||||
// as requiredIds
|
||||
// FROM query AS q
|
||||
// LEFT JOIN requiredQueries AS r ON q.id = r.queryId
|
||||
// and $2 >= r.addedVersion
|
||||
// and (r.removedVersion is null or $2 < r.removedVersion)
|
||||
// and isInVersion($2, r.addedVersion, r.removedVersion)
|
||||
// GROUP BY r.queryId
|
||||
// )
|
||||
// SELECT DISTINCT q.id, q.type, q.activeVersion, q.latestVersion, c.config,
|
||||
// SELECT DISTINCT q.id, q.type, coalesce(av.id, 0) as activeVersion, coalesce(lv.id, 0) as latestVersion, c.config,
|
||||
// coalesce(
|
||||
// r.requiredIds,
|
||||
// array[]::uuid[]
|
||||
// )::uuid[] as requiredIds
|
||||
// FROM query AS q
|
||||
// LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
// LEFT JOIN queryLatestVersions as lv on lv.queryId = q.id
|
||||
// LEFT JOIN config AS c ON q.id = c.queryId
|
||||
// LEFT JOIN requiredIds AS r ON q.id = r.queryId
|
||||
func (q *Queries) GetQueryWithVersion(ctx context.Context, arg *GetQueryWithVersionParams) (*GetQueryWithVersionRow, error) {
|
||||
@@ -411,21 +440,3 @@ func (q *Queries) RemoveRequiredQuery(ctx context.Context, arg *RemoveRequiredQu
|
||||
_, err := q.db.Exec(ctx, removeRequiredQuery, arg.Removedversion, arg.Requiredqueryid, arg.Queryid)
|
||||
return err
|
||||
}
|
||||
|
||||
const updateQuery = `-- name: UpdateQuery :exec
|
||||
UPDATE queries SET activeVersion = $1, latestVersion = $2 WHERE id = $3
|
||||
`
|
||||
|
||||
type UpdateQueryParams struct {
|
||||
Activeversion int32 `db:"activeversion"`
|
||||
Latestversion int32 `db:"latestversion"`
|
||||
ID pgtype.UUID `db:"id"`
|
||||
}
|
||||
|
||||
// UpdateQuery
|
||||
//
|
||||
// UPDATE queries SET activeVersion = $1, latestVersion = $2 WHERE id = $3
|
||||
func (q *Queries) UpdateQuery(ctx context.Context, arg *UpdateQueryParams) error {
|
||||
_, err := q.db.Exec(ctx, updateQuery, arg.Activeversion, arg.Latestversion, arg.ID)
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -36,12 +36,68 @@ func TestQueries(t *testing.T) {
|
||||
assert.NoError(t, err)
|
||||
assert.True(t, contextQueryID.Valid)
|
||||
|
||||
contextQuery, err := queries.GetQuery(ctx, contextQueryID)
|
||||
assert.NoError(t, err)
|
||||
assert.EqualExportedValues(t, &repository.Fullactivequery{
|
||||
ID: contextQueryID,
|
||||
Type: repository.QuerytypeContextFull,
|
||||
Activeversion: 0,
|
||||
Latestversion: 0,
|
||||
Config: nil,
|
||||
Requiredids: []pgtype.UUID{},
|
||||
}, contextQuery)
|
||||
|
||||
err = queries.AddLatestQueryVersion(ctx, contextQueryID)
|
||||
assert.NoError(t, err)
|
||||
|
||||
contextQuery, err = queries.GetQuery(ctx, contextQueryID)
|
||||
assert.NoError(t, err)
|
||||
assert.EqualExportedValues(t, &repository.Fullactivequery{
|
||||
ID: contextQueryID,
|
||||
Type: repository.QuerytypeContextFull,
|
||||
Activeversion: 0,
|
||||
Latestversion: 1,
|
||||
Config: nil,
|
||||
Requiredids: []pgtype.UUID{},
|
||||
}, contextQuery)
|
||||
|
||||
jsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor))
|
||||
assert.NoError(t, err)
|
||||
assert.True(t, jsonQueryID.Valid)
|
||||
|
||||
jsonQuery, err := queries.GetQuery(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
assert.EqualExportedValues(t, &repository.Fullactivequery{
|
||||
ID: jsonQueryID,
|
||||
Type: repository.QuerytypeJsonExtractor,
|
||||
Activeversion: 0,
|
||||
Latestversion: 0,
|
||||
Config: nil,
|
||||
Requiredids: []pgtype.UUID{},
|
||||
}, jsonQuery)
|
||||
|
||||
err = queries.AddLatestQueryVersion(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
|
||||
jsonQuery, err = queries.GetQuery(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
assert.EqualExportedValues(t, &repository.Fullactivequery{
|
||||
ID: jsonQueryID,
|
||||
Type: repository.QuerytypeJsonExtractor,
|
||||
Activeversion: 0,
|
||||
Latestversion: 1,
|
||||
Config: nil,
|
||||
Requiredids: []pgtype.UUID{},
|
||||
}, jsonQuery)
|
||||
|
||||
err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{
|
||||
Versionid: 1,
|
||||
Queryid: jsonQueryID,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
jsonQuery, err = queries.GetQuery(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
assert.EqualExportedValues(t, &repository.Fullactivequery{
|
||||
ID: jsonQueryID,
|
||||
Type: repository.QuerytypeJsonExtractor,
|
||||
@@ -51,13 +107,6 @@ func TestQueries(t *testing.T) {
|
||||
Requiredids: []pgtype.UUID{},
|
||||
}, jsonQuery)
|
||||
|
||||
err = queries.UpdateQuery(ctx, &repository.UpdateQueryParams{
|
||||
Latestversion: 2,
|
||||
Activeversion: 1,
|
||||
ID: jsonQueryID,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
jsonConfig := []byte("{\"path\": \"example_path\"}")
|
||||
|
||||
err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
|
||||
@@ -67,6 +116,20 @@ func TestQueries(t *testing.T) {
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
jsonQuery, err = queries.GetQuery(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
assert.EqualExportedValues(t, &repository.Fullactivequery{
|
||||
ID: jsonQueryID,
|
||||
Type: repository.QuerytypeJsonExtractor,
|
||||
Activeversion: 1,
|
||||
Latestversion: 1,
|
||||
Config: nil,
|
||||
Requiredids: []pgtype.UUID{contextQueryID},
|
||||
}, jsonQuery)
|
||||
|
||||
err = queries.AddLatestQueryVersion(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
|
||||
jsonQuery, err = queries.GetQuery(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
assert.EqualExportedValues(t, &repository.Fullactivequery{
|
||||
@@ -124,24 +187,6 @@ func TestQueries(t *testing.T) {
|
||||
Requiredids: []pgtype.UUID{contextQueryID},
|
||||
}, jsonQuery)
|
||||
|
||||
err = queries.UpdateQuery(ctx, &repository.UpdateQueryParams{
|
||||
Activeversion: 2,
|
||||
Latestversion: 2,
|
||||
ID: jsonQueryID,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
jsonQuery, err = queries.GetQuery(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
assert.EqualExportedValues(t, &repository.Fullactivequery{
|
||||
ID: jsonQueryID,
|
||||
Type: repository.QuerytypeJsonExtractor,
|
||||
Activeversion: 2,
|
||||
Latestversion: 2,
|
||||
Config: nil,
|
||||
Requiredids: []pgtype.UUID{},
|
||||
}, jsonQuery)
|
||||
|
||||
v := int32(1)
|
||||
versionedQuery, err := queries.GetQueryWithVersion(ctx, &repository.GetQueryWithVersionParams{
|
||||
ID: jsonQueryID,
|
||||
@@ -151,7 +196,7 @@ func TestQueries(t *testing.T) {
|
||||
assert.EqualExportedValues(t, &repository.GetQueryWithVersionRow{
|
||||
ID: jsonQueryID,
|
||||
Type: repository.QuerytypeJsonExtractor,
|
||||
Activeversion: 2,
|
||||
Activeversion: 1,
|
||||
Latestversion: 2,
|
||||
Config: jsonConfig,
|
||||
Requiredids: []pgtype.UUID{contextQueryID},
|
||||
@@ -219,6 +264,8 @@ func TestQueryDependencyTree(t *testing.T) {
|
||||
|
||||
jsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor))
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddLatestQueryVersion(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
|
||||
dependents, err = queries.ListQueryDirectDependentsByDocumentID(ctx, &repository.ListQueryDirectDependentsByDocumentIDParams{
|
||||
ID: docID,
|
||||
@@ -270,6 +317,8 @@ func TestQueryDependencyTree(t *testing.T) {
|
||||
|
||||
secondJsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor))
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddLatestQueryVersion(ctx, secondJsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
|
||||
dependents, err = queries.ListQueryDirectDependentsByDocumentID(ctx, &repository.ListQueryDirectDependentsByDocumentIDParams{
|
||||
ID: docID,
|
||||
@@ -403,16 +452,16 @@ func TestQueriesList(t *testing.T) {
|
||||
{
|
||||
ID: jsonQueryID,
|
||||
Type: repository.QuerytypeJsonExtractor,
|
||||
Activeversion: 1,
|
||||
Latestversion: 1,
|
||||
Activeversion: 0,
|
||||
Latestversion: 0,
|
||||
Config: nil,
|
||||
Requiredids: []pgtype.UUID{},
|
||||
},
|
||||
{
|
||||
ID: contextQueryID,
|
||||
Type: repository.QuerytypeContextFull,
|
||||
Activeversion: 1,
|
||||
Latestversion: 1,
|
||||
Activeversion: 0,
|
||||
Latestversion: 0,
|
||||
Config: nil,
|
||||
Requiredids: []pgtype.UUID{},
|
||||
},
|
||||
@@ -425,8 +474,8 @@ func TestQueriesList(t *testing.T) {
|
||||
{
|
||||
ID: jsonQueryID,
|
||||
Type: repository.QuerytypeJsonExtractor,
|
||||
Activeversion: 1,
|
||||
Latestversion: 1,
|
||||
Activeversion: 0,
|
||||
Latestversion: 0,
|
||||
Config: nil,
|
||||
Requiredids: []pgtype.UUID{},
|
||||
},
|
||||
@@ -511,6 +560,8 @@ func TestListQueryJobs(t *testing.T) {
|
||||
|
||||
jsonID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor))
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddLatestQueryVersion(ctx, jsonID)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
|
||||
Queryid: jsonID,
|
||||
Requiredqueryid: contextID,
|
||||
|
||||
@@ -46,12 +46,12 @@ func (q *Queries) GetResultValueWithVersion(ctx context.Context, arg *GetResultV
|
||||
|
||||
const listQueryRequirementValues = `-- name: ListQueryRequirementValues :many
|
||||
WITH reqQueries as (
|
||||
SELECT q.id as queryId, q.activeVersion, q.type
|
||||
SELECT q.id as queryId, av.id as activeVersion, q.type
|
||||
FROM requiredQueries as rq
|
||||
JOIN queries as q on q.id = rq.requiredQueryId
|
||||
LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
WHERE rq.queryId = $2
|
||||
and $3 >= rq.addedVersion
|
||||
and (rq.removedVersion is null or $3 < rq.removedVersion)
|
||||
and isInVersion($3, rq.addedVersion, rq.removedVersion)
|
||||
),
|
||||
docs as (
|
||||
SELECT id, jobId
|
||||
@@ -66,8 +66,7 @@ codeVersions as (
|
||||
FROM docs as d
|
||||
LEFT JOIN collectors as c on c.jobId = d.jobId
|
||||
LEFT JOIN collectorCodeVersions as ccv on c.id = ccv.collectorId
|
||||
and c.activeVersion >= ccv.addedVersion
|
||||
and c.activeVersion < COALESCE(ccv.removedVersion, c.activeVersion)
|
||||
and isInVersion(c.activeVersion, ccv.addedVersion, ccv.removedVersion)
|
||||
LIMIT 1
|
||||
),
|
||||
latestVersions AS (
|
||||
@@ -114,12 +113,12 @@ type ListQueryRequirementValuesRow struct {
|
||||
// ListQueryRequirementValues
|
||||
//
|
||||
// WITH reqQueries as (
|
||||
// SELECT q.id as queryId, q.activeVersion, q.type
|
||||
// SELECT q.id as queryId, av.id as activeVersion, q.type
|
||||
// FROM requiredQueries as rq
|
||||
// JOIN queries as q on q.id = rq.requiredQueryId
|
||||
// LEFT JOIN queryCurrentActiveVersions as av on q.id = av.queryId
|
||||
// WHERE rq.queryId = $2
|
||||
// and $3 >= rq.addedVersion
|
||||
// and (rq.removedVersion is null or $3 < rq.removedVersion)
|
||||
// and isInVersion($3, rq.addedVersion, rq.removedVersion)
|
||||
// ),
|
||||
// docs as (
|
||||
// SELECT id, jobId
|
||||
@@ -134,8 +133,7 @@ type ListQueryRequirementValuesRow struct {
|
||||
// FROM docs as d
|
||||
// LEFT JOIN collectors as c on c.jobId = d.jobId
|
||||
// LEFT JOIN collectorCodeVersions as ccv on c.id = ccv.collectorId
|
||||
// and c.activeVersion >= ccv.addedVersion
|
||||
// and c.activeVersion < COALESCE(ccv.removedVersion, c.activeVersion)
|
||||
// and isInVersion(c.activeVersion, ccv.addedVersion, ccv.removedVersion)
|
||||
// LIMIT 1
|
||||
// ),
|
||||
// latestVersions AS (
|
||||
|
||||
@@ -32,6 +32,13 @@ func TestResults(t *testing.T) {
|
||||
|
||||
jsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor))
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddLatestQueryVersion(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{
|
||||
Queryid: jsonQueryID,
|
||||
Versionid: 1,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
clientId, err := queries.CreateClient(ctx, "example_client")
|
||||
assert.NoError(t, err)
|
||||
@@ -109,6 +116,23 @@ func TestResultValues(t *testing.T) {
|
||||
contextQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeContextFull))
|
||||
assert.NoError(t, err)
|
||||
|
||||
err = queries.AddLatestQueryVersion(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{
|
||||
Queryid: jsonQueryID,
|
||||
Versionid: 1,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddLatestQueryVersion(ctx, contextQueryID)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddLatestQueryVersion(ctx, contextQueryID)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{
|
||||
Queryid: contextQueryID,
|
||||
Versionid: 2,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
jsonVersion := int32(1)
|
||||
err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
|
||||
Queryid: jsonQueryID,
|
||||
@@ -283,6 +307,20 @@ func TestUnsyncedNoDepsQueries(t *testing.T) {
|
||||
assert.NoError(t, err)
|
||||
jsonQueryID, err := queries.CreateQuery(ctx, repository.Querytype(repository.QuerytypeJsonExtractor))
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddLatestQueryVersion(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{
|
||||
Queryid: jsonQueryID,
|
||||
Versionid: 1,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddLatestQueryVersion(ctx, contextQueryID)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{
|
||||
Queryid: contextQueryID,
|
||||
Versionid: 1,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
|
||||
Queryid: jsonQueryID,
|
||||
Requiredqueryid: contextQueryID,
|
||||
@@ -326,27 +364,6 @@ func TestUnsyncedNoDepsQueries(t *testing.T) {
|
||||
assert.Len(t, qs, 1)
|
||||
assert.ElementsMatch(t, []pgtype.UUID{contextQueryID}, qs)
|
||||
|
||||
err = queries.UpdateQuery(ctx, &repository.UpdateQueryParams{
|
||||
Latestversion: 2,
|
||||
Activeversion: 2,
|
||||
ID: contextQueryID,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentID)
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 1)
|
||||
assert.ElementsMatch(t, []pgtype.UUID{contextQueryID}, qs)
|
||||
|
||||
err = queries.SetResult(ctx, &repository.SetResultParams{
|
||||
Queryid: contextQueryID,
|
||||
Documentid: documentID,
|
||||
Value: "context_value",
|
||||
Cleanversion: 1,
|
||||
Textversion: 2,
|
||||
Queryversion: 2,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
err = queries.SetResult(ctx, &repository.SetResultParams{
|
||||
Queryid: jsonQueryID,
|
||||
Documentid: documentID,
|
||||
@@ -360,11 +377,51 @@ func TestUnsyncedNoDepsQueries(t *testing.T) {
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentID)
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 0)
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentTwoID)
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 1)
|
||||
assert.ElementsMatch(t, []pgtype.UUID{contextQueryID}, qs)
|
||||
|
||||
err = queries.UpdateQuery(ctx, &repository.UpdateQueryParams{
|
||||
Latestversion: 2,
|
||||
Activeversion: 2,
|
||||
ID: jsonQueryID,
|
||||
err = queries.SetResult(ctx, &repository.SetResultParams{
|
||||
Queryid: contextQueryID,
|
||||
Documentid: documentTwoID,
|
||||
Value: "context_value",
|
||||
Cleanversion: 1,
|
||||
Textversion: 2,
|
||||
Queryversion: 1,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentID)
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 0)
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentTwoID)
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 1)
|
||||
assert.ElementsMatch(t, []pgtype.UUID{jsonQueryID}, qs)
|
||||
|
||||
err = queries.SetResult(ctx, &repository.SetResultParams{
|
||||
Queryid: jsonQueryID,
|
||||
Documentid: documentTwoID,
|
||||
Value: "context_value",
|
||||
Cleanversion: 1,
|
||||
Textversion: 2,
|
||||
Queryversion: 1,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentID)
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 0)
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentTwoID)
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 0)
|
||||
|
||||
err = queries.AddLatestQueryVersion(ctx, jsonQueryID)
|
||||
assert.NoError(t, err)
|
||||
err = queries.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{
|
||||
Queryid: jsonQueryID,
|
||||
Versionid: 2,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
@@ -372,30 +429,8 @@ func TestUnsyncedNoDepsQueries(t *testing.T) {
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 1)
|
||||
assert.ElementsMatch(t, []pgtype.UUID{jsonQueryID}, qs)
|
||||
|
||||
err = queries.SetResult(ctx, &repository.SetResultParams{
|
||||
Queryid: jsonQueryID,
|
||||
Documentid: documentID,
|
||||
Value: "context_value",
|
||||
Cleanversion: 1,
|
||||
Textversion: 2,
|
||||
Queryversion: 2,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentID)
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 0)
|
||||
|
||||
err = queries.UpdateQuery(ctx, &repository.UpdateQueryParams{
|
||||
Latestversion: 3,
|
||||
Activeversion: 3,
|
||||
ID: contextQueryID,
|
||||
})
|
||||
assert.NoError(t, err)
|
||||
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentID)
|
||||
qs, err = queries.ListUnsyncedNoDepsQueriesByDocId(ctx, documentTwoID)
|
||||
assert.NoError(t, err)
|
||||
assert.Len(t, qs, 1)
|
||||
assert.ElementsMatch(t, []pgtype.UUID{contextQueryID}, qs)
|
||||
assert.ElementsMatch(t, []pgtype.UUID{jsonQueryID}, qs)
|
||||
}
|
||||
|
||||
@@ -64,12 +64,26 @@ func (s *Service) submitCreate(ctx context.Context, entity *resultprocessor.Crea
|
||||
return err
|
||||
}
|
||||
|
||||
err = qtx.AddLatestQueryVersion(ctx, dbID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
version := int32(1)
|
||||
|
||||
err = qtx.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{
|
||||
Queryid: dbID,
|
||||
Versionid: version,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if query.RequiredQueryIDs != nil {
|
||||
for _, reqQuery := range *query.RequiredQueryIDs {
|
||||
err = qtx.AddRequiredQuery(ctx, &repository.AddRequiredQueryParams{
|
||||
Queryid: dbID,
|
||||
Requiredqueryid: reqQuery,
|
||||
Addedversion: 1,
|
||||
Addedversion: version,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -77,11 +91,11 @@ func (s *Service) submitCreate(ctx context.Context, entity *resultprocessor.Crea
|
||||
}
|
||||
}
|
||||
|
||||
if query.Config != nil && string(*query.Config) != "" {
|
||||
if query.Config != nil {
|
||||
err = qtx.AddQueryConfig(ctx, &repository.AddQueryConfigParams{
|
||||
Queryid: dbID,
|
||||
Config: *query.Config,
|
||||
Addedversion: 1,
|
||||
Addedversion: version,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
@@ -55,6 +55,10 @@ func TestCreate(t *testing.T) {
|
||||
pgxmock.NewRows([]string{"id"}).
|
||||
AddRow(database.MustToDBUUID(q.ID)),
|
||||
)
|
||||
pool.ExpectExec("name: AddLatestQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID), int32(1)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
for _, req := range *create.RequiredQueryIDs {
|
||||
pool.ExpectExec("name: AddRequiredQuery :exec").WithArgs(database.MustToDBUUID(q.ID), database.MustToDBUUID(req), int32(1)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
@@ -96,6 +100,10 @@ func TestCreateMinimal(t *testing.T) {
|
||||
pgxmock.NewRows([]string{"id"}).
|
||||
AddRow(database.MustToDBUUID(q.ID)),
|
||||
)
|
||||
pool.ExpectExec("name: AddLatestQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID), int32(1)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectCommit()
|
||||
|
||||
id, err := svc.Create(ctx, create)
|
||||
|
||||
@@ -109,6 +109,10 @@ func TestSubmitCreate(t *testing.T) {
|
||||
pgxmock.NewRows([]string{"id"}).
|
||||
AddRow(database.MustToDBUUID(q.ID)),
|
||||
)
|
||||
pool.ExpectExec("name: AddLatestQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID), int32(1)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
for _, req := range *create.RequiredQueryIDs {
|
||||
pool.ExpectExec("name: AddRequiredQuery :exec").WithArgs(database.MustToDBUUID(q.ID), database.MustToDBUUID(req), int32(1)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
@@ -150,6 +154,10 @@ func TestSubmitCreateNoReqsOrConfig(t *testing.T) {
|
||||
pgxmock.NewRows([]string{"id"}).
|
||||
AddRow(database.MustToDBUUID(q.ID)),
|
||||
)
|
||||
pool.ExpectExec("name: AddLatestQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID), int32(1)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectCommit()
|
||||
|
||||
id, err := svc.submitCreate(ctx, create)
|
||||
|
||||
@@ -126,9 +126,33 @@ func (s *Service) normalizeUpdate(ctx context.Context, current *query.Query, ent
|
||||
|
||||
func (s *Service) submitUpdate(ctx context.Context, current *query.Query, entity *resultprocessor.Update) error {
|
||||
err := s.cfg.ExecuteDBTransaction(ctx, func(ctx context.Context, q *repository.Queries) error {
|
||||
latestVersion := current.LatestVersion + 1
|
||||
id := database.MustToDBUUID(entity.ID)
|
||||
|
||||
activeName, err := validation.GetFieldName(entity, entity.ActiveVersion)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
onlyactive := validation.AreAllPointersNilExcept(entity, activeName)
|
||||
|
||||
if !onlyactive {
|
||||
err := q.AddLatestQueryVersion(ctx, id)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
if entity.ActiveVersion != nil {
|
||||
err := q.AddActiveQueryVersion(ctx, &repository.AddActiveQueryVersionParams{
|
||||
Queryid: id,
|
||||
Versionid: *entity.ActiveVersion,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
latestVersion := current.LatestVersion + 1
|
||||
|
||||
if entity.RequiredQueryIDs != nil {
|
||||
addIDs := getSetDifference(entity.RequiredQueryIDs, current.RequiredQueryIDs)
|
||||
for _, qID := range addIDs {
|
||||
@@ -174,30 +198,6 @@ func (s *Service) submitUpdate(ctx context.Context, current *query.Query, entity
|
||||
}
|
||||
}
|
||||
|
||||
activeVersion := entity.ActiveVersion
|
||||
if activeVersion == nil {
|
||||
activeVersion = ¤t.ActiveVersion
|
||||
}
|
||||
|
||||
activeName, err := validation.GetFieldName(entity, entity.ActiveVersion)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
onlyactive := validation.AreAllPointersNilExcept(entity, activeName)
|
||||
|
||||
if onlyactive {
|
||||
latestVersion--
|
||||
}
|
||||
|
||||
err = q.UpdateQuery(ctx, &repository.UpdateQueryParams{
|
||||
Latestversion: latestVersion,
|
||||
Activeversion: *activeVersion,
|
||||
ID: id,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
slog.Debug("query updated", "update", *entity)
|
||||
|
||||
return nil
|
||||
|
||||
@@ -63,12 +63,14 @@ func TestUpdate(t *testing.T) {
|
||||
AddRow(database.MustToDBUUID(existing.ID), repository.QuerytypeJsonExtractor, existing.ActiveVersion, existing.LatestVersion, []byte(config), []pgtype.UUID{}),
|
||||
)
|
||||
pool.ExpectBeginTx(pgx.TxOptions{})
|
||||
pool.ExpectExec("name: AddLatestQueryVersion :exec").WithArgs(database.MustToDBUUID(update.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(update.ID), av).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: RemoveQueryConfig :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID(update.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddQueryConfig :exec").WithArgs(database.MustToDBUUID(update.ID), []byte(*update.Config), int32(2)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(av, int32(2), database.MustToDBUUID(update.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectCommit()
|
||||
|
||||
mockSQS.EXPECT().
|
||||
|
||||
@@ -88,6 +88,10 @@ func TestSubmitUpdate(t *testing.T) {
|
||||
}
|
||||
|
||||
pool.ExpectBeginTx(pgx.TxOptions{})
|
||||
pool.ExpectExec("name: AddLatestQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID), aV).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddRequiredQuery :exec").WithArgs(database.MustToDBUUID(update.ID), database.MustToDBUUID((*update.RequiredQueryIDs)[0]), pgxmock.AnyArg()).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: RemoveRequiredQuery :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID((*q.RequiredQueryIDs)[0]), database.MustToDBUUID(update.ID)).
|
||||
@@ -96,8 +100,6 @@ func TestSubmitUpdate(t *testing.T) {
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddQueryConfig :exec").WithArgs(database.MustToDBUUID(update.ID), []byte(*update.Config), int32(3)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(aV, int32(3), database.MustToDBUUID(update.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectCommit()
|
||||
|
||||
err = svc.submitUpdate(ctx, &q, update)
|
||||
@@ -131,7 +133,7 @@ func TestSubmitUpdateRollback(t *testing.T) {
|
||||
|
||||
pool.ExpectBeginTx(pgx.TxOptions{})
|
||||
msg := "database failure"
|
||||
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(aV, int32(2), database.MustToDBUUID(update.ID)).
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID), aV).
|
||||
WillReturnError(errors.New(msg))
|
||||
pool.ExpectCommit()
|
||||
|
||||
@@ -177,6 +179,8 @@ func TestSubmitUpdateRequiredQueries(t *testing.T) {
|
||||
}
|
||||
|
||||
pool.ExpectBeginTx(pgx.TxOptions{})
|
||||
pool.ExpectExec("name: AddLatestQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddRequiredQuery :exec").WithArgs(database.MustToDBUUID(update.ID), database.MustToDBUUID((*update.RequiredQueryIDs)[0]), pgxmock.AnyArg()).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: AddRequiredQuery :exec").WithArgs(database.MustToDBUUID(update.ID), database.MustToDBUUID((*update.RequiredQueryIDs)[1]), pgxmock.AnyArg()).
|
||||
@@ -185,8 +189,6 @@ func TestSubmitUpdateRequiredQueries(t *testing.T) {
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: RemoveRequiredQuery :exec").WithArgs(pgxmock.AnyArg(), database.MustToDBUUID((*q.RequiredQueryIDs)[2]), database.MustToDBUUID(update.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(q.ActiveVersion, int32(3), database.MustToDBUUID(update.ID)).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectCommit()
|
||||
|
||||
err = svc.submitUpdate(ctx, &q, update)
|
||||
@@ -226,7 +228,7 @@ func TestSubmitUpdateRequiredQueries(t *testing.T) {
|
||||
}
|
||||
|
||||
pool.ExpectBeginTx(pgx.TxOptions{})
|
||||
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(av, int32(2), database.MustToDBUUID(update.ID)).
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID), av).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectCommit()
|
||||
|
||||
@@ -261,7 +263,7 @@ func TestSubmitUpdateActiveVersion(t *testing.T) {
|
||||
}
|
||||
|
||||
pool.ExpectBeginTx(pgx.TxOptions{})
|
||||
pool.ExpectExec("name: UpdateQuery :exec").WithArgs(aV, int32(2), database.MustToDBUUID(update.ID)).
|
||||
pool.ExpectExec("name: AddActiveQueryVersion :exec").WithArgs(database.MustToDBUUID(q.ID), aV).
|
||||
WillReturnResult(pgxmock.NewResult("", 1))
|
||||
pool.ExpectCommit()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user