feat: added initial version pruning and result limiting

This commit is contained in:
Sheen Capadngan
2024-05-31 19:12:55 +08:00
parent b801c1e48f
commit 3a1168c7e8
8 changed files with 140 additions and 7 deletions
@@ -0,0 +1,21 @@
import { Knex } from "knex";
import { TableName } from "../schemas";
export async function up(knex: Knex): Promise<void> {
const hasPitVersionLimitColumn = await knex.schema.hasColumn(TableName.Project, "pitVersionLimit");
await knex.schema.alterTable(TableName.Project, (tb) => {
if (!hasPitVersionLimitColumn) {
tb.integer("pitVersionLimit").notNullable().defaultTo(10);
}
});
}
export async function down(knex: Knex): Promise<void> {
const hasPitVersionLimitColumn = await knex.schema.hasColumn(TableName.Project, "pitVersionLimit");
await knex.schema.alterTable(TableName.Project, (tb) => {
if (hasPitVersionLimitColumn) {
tb.dropColumn("pitVersionLimit");
}
});
}
+2 -1
View File
@@ -16,7 +16,8 @@ export const ProjectsSchema = z.object({
createdAt: z.date(), createdAt: z.date(),
updatedAt: z.date(), updatedAt: z.date(),
version: z.number().default(1), version: z.number().default(1),
upgradeStatus: z.string().nullable().optional() upgradeStatus: z.string().nullable().optional(),
pitVersionLimit: z.number()
}); });
export type TProjects = z.infer<typeof ProjectsSchema>; export type TProjects = z.infer<typeof ProjectsSchema>;
@@ -4,6 +4,7 @@ import { TableName, TSecretTagJunctionInsert } from "@app/db/schemas";
import { BadRequestError, InternalServerError } from "@app/lib/errors"; import { BadRequestError, InternalServerError } from "@app/lib/errors";
import { groupBy } from "@app/lib/fn"; import { groupBy } from "@app/lib/fn";
import { logger } from "@app/lib/logger"; import { logger } from "@app/lib/logger";
import { TProjectDALFactory } from "@app/services/project/project-dal";
import { TSecretDALFactory } from "@app/services/secret/secret-dal"; import { TSecretDALFactory } from "@app/services/secret/secret-dal";
import { TSecretVersionDALFactory } from "@app/services/secret/secret-version-dal"; import { TSecretVersionDALFactory } from "@app/services/secret/secret-version-dal";
import { TSecretVersionTagDALFactory } from "@app/services/secret/secret-version-tag-dal"; import { TSecretVersionTagDALFactory } from "@app/services/secret/secret-version-tag-dal";
@@ -37,6 +38,7 @@ type TSecretSnapshotServiceFactoryDep = {
folderDAL: Pick<TSecretFolderDALFactory, "findById" | "findBySecretPath" | "delete" | "insertMany" | "find">; folderDAL: Pick<TSecretFolderDALFactory, "findById" | "findBySecretPath" | "delete" | "insertMany" | "find">;
permissionService: Pick<TPermissionServiceFactory, "getProjectPermission">; permissionService: Pick<TPermissionServiceFactory, "getProjectPermission">;
licenseService: Pick<TLicenseServiceFactory, "isValidLicense">; licenseService: Pick<TLicenseServiceFactory, "isValidLicense">;
projectDAL: Pick<TProjectDALFactory, "findById">;
}; };
export type TSecretSnapshotServiceFactory = ReturnType<typeof secretSnapshotServiceFactory>; export type TSecretSnapshotServiceFactory = ReturnType<typeof secretSnapshotServiceFactory>;
@@ -48,6 +50,7 @@ export const secretSnapshotServiceFactory = ({
snapshotSecretDAL, snapshotSecretDAL,
snapshotFolderDAL, snapshotFolderDAL,
folderDAL, folderDAL,
projectDAL,
secretDAL, secretDAL,
permissionService, permissionService,
licenseService, licenseService,
@@ -81,8 +84,9 @@ export const secretSnapshotServiceFactory = ({
const folder = await folderDAL.findBySecretPath(projectId, environment, path); const folder = await folderDAL.findBySecretPath(projectId, environment, path);
if (!folder) throw new BadRequestError({ message: "Folder not found" }); if (!folder) throw new BadRequestError({ message: "Folder not found" });
const project = await projectDAL.findById(projectId);
const count = await snapshotDAL.countOfSnapshotsByFolderId(folder.id); const count = await snapshotDAL.countOfSnapshotsByFolderId(folder.id);
return count; return Math.min(count, project.pitVersionLimit);
}; };
const listSnapshots = async ({ const listSnapshots = async ({
@@ -114,7 +118,16 @@ export const secretSnapshotServiceFactory = ({
const folder = await folderDAL.findBySecretPath(projectId, environment, path); const folder = await folderDAL.findBySecretPath(projectId, environment, path);
if (!folder) throw new BadRequestError({ message: "Folder not found" }); if (!folder) throw new BadRequestError({ message: "Folder not found" });
const snapshots = await snapshotDAL.find({ folderId: folder.id }, { limit, offset, sort: [["createdAt", "desc"]] }); const { pitVersionLimit } = await projectDAL.findById(projectId);
const computedQueryLimit = Math.min(pitVersionLimit - offset, limit);
if (offset > pitVersionLimit || computedQueryLimit <= 0) {
return [];
}
const snapshots = await snapshotDAL.find(
{ folderId: folder.id },
{ limit: computedQueryLimit, offset, sort: [["createdAt", "desc"]] }
);
return snapshots; return snapshots;
}; };
@@ -325,12 +325,62 @@ export const snapshotDALFactory = (db: TDbClient) => {
} }
}; };
const pruneExcessSnapshots = async (tx?: Knex) => {
try {
const folders = await (tx || db)(TableName.SecretFolder).select("id");
const folderIds = folders.map((folder) => folder.id);
const PRUNE_FOLDER_BATCH_SIZE = 500;
const pruneBatches = [];
for (let x = 0; x < folderIds.length; x += PRUNE_FOLDER_BATCH_SIZE) {
const batch = folderIds.slice(x, x + PRUNE_FOLDER_BATCH_SIZE);
pruneBatches.push(batch);
}
for await (const folderBatch of pruneBatches) {
const rankedSnapshots = (tx || db)(TableName.Snapshot)
.whereIn(`${TableName.Snapshot}.folderId`, folderBatch)
.select(
"folderId",
"id",
(tx || db).raw(
`ROW_NUMBER() OVER (PARTITION BY ${TableName.Snapshot}."folderId" ORDER BY ${TableName.Snapshot}."createdAt" DESC) AS row_num`
)
)
.as("ranked_snapshots");
const snapshotsToKeep = (tx || db)
.select("id")
.from(rankedSnapshots)
.where(
"row_num",
"<=",
(tx || db)
.select(`${TableName.Project}.pitVersionLimit`)
.from(TableName.Project)
.join(TableName.Environment, `${TableName.Environment}.projectId`, `${TableName.Project}.id`)
.join(TableName.Snapshot, `${TableName.Snapshot}.envId`, `${TableName.Environment}.id`)
.join(rankedSnapshots, "ranked_snapshots.folderId", `${TableName.Snapshot}.folderId`)
.limit(1)
);
await (tx || db)(TableName.Snapshot)
.whereIn("folderId", folderBatch)
.whereNotIn("id", snapshotsToKeep)
.delete();
}
} catch (error) {
throw new DatabaseError({ error, name: "SnapshotPrune" });
}
};
return { return {
...secretSnapshotOrm, ...secretSnapshotOrm,
findById, findById,
findLatestSnapshotByFolderId, findLatestSnapshotByFolderId,
findRecursivelySnapshots, findRecursivelySnapshots,
countOfSnapshotsByFolderId, countOfSnapshotsByFolderId,
findSecretSnapshotDataById findSecretSnapshotDataById,
pruneExcessSnapshots
}; };
}; };
+1
View File
@@ -535,6 +535,7 @@ export const registerRoutes = async (
licenseService, licenseService,
folderDAL, folderDAL,
secretDAL, secretDAL,
projectDAL,
snapshotDAL, snapshotDAL,
snapshotFolderDAL, snapshotFolderDAL,
snapshotSecretDAL, snapshotSecretDAL,
@@ -133,7 +133,8 @@ export const projectServiceFactory = ({
name: workspaceName, name: workspaceName,
orgId: organization.id, orgId: organization.id,
slug: projectSlug || slugify(`${workspaceName}-${alphaNumericNanoId(4)}`), slug: projectSlug || slugify(`${workspaceName}-${alphaNumericNanoId(4)}`),
version: ProjectVersion.V2 version: ProjectVersion.V2,
pitVersionLimit: 10
}, },
tx tx
); );
+11 -2
View File
@@ -72,7 +72,7 @@ type TSecretServiceFactoryDep = {
secretDAL: TSecretDALFactory; secretDAL: TSecretDALFactory;
secretTagDAL: TSecretTagDALFactory; secretTagDAL: TSecretTagDALFactory;
secretVersionDAL: TSecretVersionDALFactory; secretVersionDAL: TSecretVersionDALFactory;
projectDAL: Pick<TProjectDALFactory, "checkProjectUpgradeStatus" | "findProjectBySlug">; projectDAL: Pick<TProjectDALFactory, "checkProjectUpgradeStatus" | "findProjectBySlug" | "findById">;
projectEnvDAL: Pick<TProjectEnvDALFactory, "findOne">; projectEnvDAL: Pick<TProjectEnvDALFactory, "findOne">;
folderDAL: Pick< folderDAL: Pick<
TSecretFolderDALFactory, TSecretFolderDALFactory,
@@ -1354,7 +1354,16 @@ export const secretServiceFactory = ({
); );
ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionActions.Read, ProjectPermissionSub.SecretRollback); ForbiddenError.from(permission).throwUnlessCan(ProjectPermissionActions.Read, ProjectPermissionSub.SecretRollback);
const secretVersions = await secretVersionDAL.find({ secretId }, { offset, limit, sort: [["createdAt", "desc"]] }); const { pitVersionLimit } = await projectDAL.findById(folder.projectId);
const computedQueryLimit = Math.min(pitVersionLimit - offset, limit);
if (offset > pitVersionLimit || computedQueryLimit <= 0) {
return [];
}
const secretVersions = await secretVersionDAL.find(
{ secretId },
{ offset, limit: computedQueryLimit, sort: [["createdAt", "desc"]] }
);
return secretVersions; return secretVersions;
}; };
@@ -110,8 +110,45 @@ export const secretVersionDALFactory = (db: TDbClient) => {
} }
}; };
const pruneExcessVersions = async (tx?: Knex) => {
try {
const rankedSecretVersions = (tx || db)(TableName.SecretVersion)
.select(
"id",
"secretId",
"folderId",
(tx || db).raw(
`ROW_NUMBER() OVER (PARTITION BY ${TableName.SecretVersion}."secretId" ORDER BY ${TableName.SecretVersion}."createdAt" DESC) AS row_num`
)
)
.as("ranked_secret_versions");
const versionsToKeep = (tx || db)(rankedSecretVersions)
.select("id")
.where(
"row_num",
"<=",
(tx || db)
.select(`${TableName.Project}.pitVersionLimit`)
.from(TableName.Project)
.join(TableName.Environment, `${TableName.Environment}.projectId`, `${TableName.Project}.id`)
.join(TableName.SecretFolder, `${TableName.SecretFolder}.envId`, `${TableName.Environment}.id`)
.join(rankedSecretVersions, "ranked_secret_versions.folderId", `${TableName.SecretFolder}.id`)
.limit(1)
);
await (tx || db)(TableName.SecretVersion).whereNotIn("id", versionsToKeep).delete();
} catch (error) {
throw new DatabaseError({
error,
name: "Secret Version Prune"
});
}
};
return { return {
...secretVersionOrm, ...secretVersionOrm,
pruneExcessVersions,
findLatestVersionMany, findLatestVersionMany,
bulkUpdate, bulkUpdate,
findLatestVersionByFolderId, findLatestVersionByFolderId,