From 8bfd3913da82c95923d6e52543fa35b3e7805a83 Mon Sep 17 00:00:00 2001 From: carlosmonastyrski Date: Wed, 14 May 2025 10:26:41 -0300 Subject: [PATCH] PIT: add backend logic for deep PIT and rollback --- .../20250505194916_add-pit-revamp-tables.ts | 5 +- backend/src/db/migrations/utils/services.ts | 5 +- backend/src/db/schemas/folder-commits.ts | 1 + backend/src/ee/routes/v1/pit-router.ts | 109 +- backend/src/lib/config/env.ts | 3 +- backend/src/queue/queue-service.ts | 12 +- backend/src/server/routes/index.ts | 14 +- .../folder-commit/folder-commit-dal.ts | 157 ++- .../folder-commit/folder-commit-queue.ts | 180 ++++ .../folder-commit-service.test.ts | 571 +++++++++++ .../folder-commit/folder-commit-service.ts | 952 ++++++++++++------ .../folder-tree-checkpoint-resources-dal.ts | 110 +- .../folder-tree-checkpoint-dal.ts | 101 +- .../secret-folder/secret-folder-dal.ts | 60 +- 14 files changed, 1818 insertions(+), 462 deletions(-) create mode 100644 backend/src/services/folder-commit/folder-commit-queue.ts create mode 100644 backend/src/services/folder-commit/folder-commit-service.test.ts diff --git a/backend/src/db/migrations/20250505194916_add-pit-revamp-tables.ts b/backend/src/db/migrations/20250505194916_add-pit-revamp-tables.ts index 9855d301e..a2a8ac8c2 100644 --- a/backend/src/db/migrations/20250505194916_add-pit-revamp-tables.ts +++ b/backend/src/db/migrations/20250505194916_add-pit-revamp-tables.ts @@ -13,10 +13,12 @@ export async function up(knex: Knex): Promise { t.string("actorType").notNullable(); t.string("message"); t.uuid("folderId").notNullable(); - t.foreign("folderId").references("id").inTable(TableName.SecretFolder).onDelete("CASCADE"); + t.uuid("envId").notNullable(); + t.foreign("envId").references("id").inTable(TableName.Environment).onDelete("CASCADE"); t.timestamps(true, true, true); t.index("folderId"); + t.index("envId"); }); } @@ -88,7 +90,6 @@ export async function up(knex: Knex): Promise { t.uuid("folderTreeCheckpointId").notNullable(); t.foreign("folderTreeCheckpointId").references("id").inTable(TableName.FolderTreeCheckpoint).onDelete("CASCADE"); t.uuid("folderId").notNullable(); - t.foreign("folderId").references("id").inTable(TableName.SecretFolder).onDelete("CASCADE"); t.uuid("folderCommitId").notNullable(); t.foreign("folderCommitId").references("id").inTable(TableName.FolderCommit).onDelete("CASCADE"); t.timestamps(true, true, true); diff --git a/backend/src/db/migrations/utils/services.ts b/backend/src/db/migrations/utils/services.ts index 1ee81536c..4e0188c98 100644 --- a/backend/src/db/migrations/utils/services.ts +++ b/backend/src/db/migrations/utils/services.ts @@ -9,6 +9,7 @@ import { folderCommitDALFactory } from "@app/services/folder-commit/folder-commi import { folderCommitServiceFactory } from "@app/services/folder-commit/folder-commit-service"; import { folderCommitChangesDALFactory } from "@app/services/folder-commit-changes/folder-commit-changes-dal"; import { folderTreeCheckpointDALFactory } from "@app/services/folder-tree-checkpoint/folder-tree-checkpoint-dal"; +import { folderTreeCheckpointResourcesDALFactory } from "@app/services/folder-tree-checkpoint-resources/folder-tree-checkpoint-resources-dal"; import { identityDALFactory } from "@app/services/identity/identity-dal"; import { internalKmsDALFactory } from "@app/services/kms/internal-kms-dal"; import { kmskeyDALFactory } from "@app/services/kms/kms-key-dal"; @@ -76,6 +77,7 @@ export const getMigrationPITServices = async ({ db, keyStore }: { db: Knex; keyS const secretVersionV2BridgeDAL = secretVersionV2BridgeDALFactory(db); const folderCheckpointResourcesDAL = folderCheckpointResourcesDALFactory(db); const secretV2BridgeDAL = secretV2BridgeDALFactory({ db, keyStore }); + const folderTreeCheckpointResourcesDAL = folderTreeCheckpointResourcesDALFactory(db); const folderCommitService = folderCommitServiceFactory({ folderCommitDAL, @@ -89,7 +91,8 @@ export const getMigrationPITServices = async ({ db, keyStore }: { db: Knex; keyS secretVersionV2BridgeDAL, projectDAL, folderCheckpointResourcesDAL, - secretV2BridgeDAL + secretV2BridgeDAL, + folderTreeCheckpointResourcesDAL }); return { folderCommitService }; diff --git a/backend/src/db/schemas/folder-commits.ts b/backend/src/db/schemas/folder-commits.ts index ec166e7c2..fc9b83489 100644 --- a/backend/src/db/schemas/folder-commits.ts +++ b/backend/src/db/schemas/folder-commits.ts @@ -14,6 +14,7 @@ export const FolderCommitsSchema = z.object({ actorType: z.string(), message: z.string().nullable().optional(), folderId: z.string().uuid(), + envId: z.string().uuid(), createdAt: z.date(), updatedAt: z.date() }); diff --git a/backend/src/ee/routes/v1/pit-router.ts b/backend/src/ee/routes/v1/pit-router.ts index 9abef5382..a3ec8dcc9 100644 --- a/backend/src/ee/routes/v1/pit-router.ts +++ b/backend/src/ee/routes/v1/pit-router.ts @@ -3,9 +3,52 @@ import { z } from "zod"; import { readLimit } from "@app/server/config/rateLimiter"; export const registerPITRouter = async (server: FastifyZodProvider) => { + // Get all commits for a folder server.route({ method: "GET", - url: "/diff", + url: "/commits/:folderId", + config: { + rateLimit: readLimit + }, + schema: { + params: z.object({ + folderId: z.string().trim() + }), + response: { + 200: z.any() + } + }, + handler: async (req) => { + const commits = await server.services.folderCommit.getCommitsByFolderId(req.params.folderId); + return commits; + } + }); + + // Get commit changes for a specific commit + server.route({ + method: "GET", + url: "/commits/:commitId/changes", + config: { + rateLimit: readLimit + }, + schema: { + params: z.object({ + commitId: z.string().trim() + }), + response: { + 200: z.any() + } + }, + handler: async (req) => { + const changes = await server.services.folderCommit.getCommitChanges(req.params.commitId); + return changes; + } + }); + + // Compare folder states between commits + server.route({ + method: "GET", + url: "/compare", config: { rateLimit: readLimit }, @@ -19,12 +62,16 @@ export const registerPITRouter = async (server: FastifyZodProvider) => { } }, handler: async (req) => { - const backup = await server.services.folderCommit.compareFolderStates(req.query.fromCommit, req.query.toCommit); + const diff = await server.services.folderCommit.compareFolderStates({ + currentCommitId: req.query.fromCommit, + targetCommitId: req.query.toCommit + }); - return backup; + return diff; } }); + // Rollback to a previous commit server.route({ method: "POST", url: "/rollback", @@ -36,6 +83,46 @@ export const registerPITRouter = async (server: FastifyZodProvider) => { fromCommit: z.string().trim(), toCommit: z.string().trim(), folderId: z.string().trim(), + projectId: z.string().trim(), + reconstructNewFolders: z.boolean().default(false) + }), + response: { + 200: z.any() + } + }, + handler: async (req) => { + const diff = await server.services.folderCommit.compareFolderStates({ + currentCommitId: req.body.fromCommit, + targetCommitId: req.body.toCommit + }); + + const response = await server.services.folderCommit.applyFolderStateDifferences({ + differences: diff, + actorInfo: { + actorType: req.permission?.type || "PLATFORM", + actorId: req.permission?.id, + message: "Rollback to previous commit" + }, + folderId: req.body.folderId, + projectId: req.body.projectId, + reconstructNewFolders: req.body.reconstructNewFolders + }); + + return response; + } + }); + + // Deep rollback to a specific commit + server.route({ + method: "POST", + url: "/deep-rollback", + config: { + rateLimit: readLimit + }, + schema: { + body: z.object({ + commitId: z.string().trim(), + envId: z.string().trim(), projectId: z.string().trim() }), response: { @@ -43,19 +130,15 @@ export const registerPITRouter = async (server: FastifyZodProvider) => { } }, handler: async (req) => { - const diff = await server.services.folderCommit.compareFolderStates(req.body.fromCommit, req.body.toCommit); - const response = await server.services.folderCommit.applyFolderStateDifferences( - diff, - { - actorType: req.permission?.type || "PLATFORM", - actorId: req.permission?.id, - message: "Rollback to previous commit" - }, - req.body.folderId, + await server.services.folderCommit.deepRollbackFolder( + req.body.commitId, + req.body.envId, + req.permission?.id || "PLATFORM", + req.permission?.type || "PLATFORM", req.body.projectId ); - return response; + return { success: true }; } }); }; diff --git a/backend/src/lib/config/env.ts b/backend/src/lib/config/env.ts index 38e252b2a..74da0b399 100644 --- a/backend/src/lib/config/env.ts +++ b/backend/src/lib/config/env.ts @@ -229,7 +229,8 @@ const envSchema = z DATADOG_HOSTNAME: zpStr(z.string().optional()), // PIT - CHECKPOINT_WINDOW: zpStr(z.string().optional().default("10")), + PIT_CHECKPOINT_WINDOW: zpStr(z.string().optional().default("10")), + PIT_TREE_CHECKPOINT_WINDOW: zpStr(z.string().optional().default("100")), /* CORS ----------------------------------------------------------------------------- */ diff --git a/backend/src/queue/queue-service.ts b/backend/src/queue/queue-service.ts index ae1a3e821..f02519e3a 100644 --- a/backend/src/queue/queue-service.ts +++ b/backend/src/queue/queue-service.ts @@ -49,7 +49,8 @@ export enum QueueName { AccessTokenStatusUpdate = "access-token-status-update", ImportSecretsFromExternalSource = "import-secrets-from-external-source", AppConnectionSecretSync = "app-connection-secret-sync", - SecretRotationV2 = "secret-rotation-v2" + SecretRotationV2 = "secret-rotation-v2", + FolderTreeCheckpoint = "folder-tree-checkpoint" } export enum QueueJobs { @@ -81,7 +82,8 @@ export enum QueueJobs { SecretSyncSendActionFailedNotifications = "secret-sync-send-action-failed-notifications", SecretRotationV2QueueRotations = "secret-rotation-v2-queue-rotations", SecretRotationV2RotateSecrets = "secret-rotation-v2-rotate-secrets", - SecretRotationV2SendNotification = "secret-rotation-v2-send-notification" + SecretRotationV2SendNotification = "secret-rotation-v2-send-notification", + CreateFolderTreeCheckpoint = "create-folder-tree-checkpoint" } export type TQueueJobTypes = { @@ -191,6 +193,12 @@ export type TQueueJobTypes = { name: QueueJobs.ProjectV3Migration; payload: { projectId: string }; }; + [QueueName.FolderTreeCheckpoint]: { + name: QueueJobs.CreateFolderTreeCheckpoint; + payload: { + envId: string; + }; + }; [QueueName.ImportSecretsFromExternalSource]: { name: QueueJobs.ImportSecretsFromExternalSource; payload: { diff --git a/backend/src/server/routes/index.ts b/backend/src/server/routes/index.ts index 38c5f5dd0..2960519bd 100644 --- a/backend/src/server/routes/index.ts +++ b/backend/src/server/routes/index.ts @@ -143,9 +143,11 @@ import { externalMigrationServiceFactory } from "@app/services/external-migratio import { folderCheckpointDALFactory } from "@app/services/folder-checkpoint/folder-checkpoint-dal"; import { folderCheckpointResourcesDALFactory } from "@app/services/folder-checkpoint-resources/folder-checkpoint-resources-dal"; import { folderCommitDALFactory } from "@app/services/folder-commit/folder-commit-dal"; +import { folderCommitQueueServiceFactory } from "@app/services/folder-commit/folder-commit-queue"; import { folderCommitServiceFactory } from "@app/services/folder-commit/folder-commit-service"; import { folderCommitChangesDALFactory } from "@app/services/folder-commit-changes/folder-commit-changes-dal"; import { folderTreeCheckpointDALFactory } from "@app/services/folder-tree-checkpoint/folder-tree-checkpoint-dal"; +import { folderTreeCheckpointResourcesDALFactory } from "@app/services/folder-tree-checkpoint-resources/folder-tree-checkpoint-resources-dal"; import { groupProjectDALFactory } from "@app/services/group-project/group-project-dal"; import { groupProjectMembershipRoleDALFactory } from "@app/services/group-project/group-project-membership-role-dal"; import { groupProjectServiceFactory } from "@app/services/group-project/group-project-service"; @@ -563,6 +565,14 @@ export const registerRoutes = async ( const folderCheckpointResourcesDAL = folderCheckpointResourcesDALFactory(db); const folderTreeCheckpointDAL = folderTreeCheckpointDALFactory(db); const folderCommitDAL = folderCommitDALFactory(db); + const folderTreeCheckpointResourcesDAL = folderTreeCheckpointResourcesDALFactory(db); + const folderCommitQueueService = folderCommitQueueServiceFactory({ + queueService, + folderTreeCheckpointDAL, + folderTreeCheckpointResourcesDAL, + folderCommitDAL, + folderDAL + }); const folderCommitService = folderCommitServiceFactory({ folderCommitDAL, folderCommitChangesDAL, @@ -575,7 +585,9 @@ export const registerRoutes = async ( secretVersionV2BridgeDAL, projectDAL, folderCheckpointResourcesDAL, - secretV2BridgeDAL + secretV2BridgeDAL, + folderTreeCheckpointResourcesDAL, + folderCommitQueueService }); const scimService = scimServiceFactory({ licenseService, diff --git a/backend/src/services/folder-commit/folder-commit-dal.ts b/backend/src/services/folder-commit/folder-commit-dal.ts index f52da95b5..7f015ddc0 100644 --- a/backend/src/services/folder-commit/folder-commit-dal.ts +++ b/backend/src/services/folder-commit/folder-commit-dal.ts @@ -42,6 +42,75 @@ export const folderCommitDALFactory = (db: TDbClient) => { } }; + const findLatestCommitByFolderIds = async (folderIds: string[], tx?: Knex): Promise => { + try { + // First get max commitId for each folderId + const maxCommitIdSubquery = (tx || db.replicaNode())(TableName.FolderCommit) + .select("folderId") + .max("commitId as maxCommitId") + .whereIn("folderId", folderIds) + .groupBy("folderId"); + + // Join with main table to get complete records for each max commitId + const docs = await (tx || db.replicaNode())(TableName.FolderCommit) + .select(selectAllTableCols(TableName.FolderCommit)) + // eslint-disable-next-line func-names + .join(maxCommitIdSubquery.as("latest"), function () { + this.on(`${TableName.FolderCommit}.folderId`, "=", "latest.folderId").andOn( + `${TableName.FolderCommit}.commitId`, + "=", + "latest.maxCommitId" + ); + }); + + return docs; + } catch (error) { + throw new DatabaseError({ error, name: "FindLatestCommitByFolderIds" }); + } + }; + + const findLatestEnvCommit = async (envId: string, tx?: Knex): Promise => { + try { + const doc = await (tx || db.replicaNode())(TableName.FolderCommit) + .where(`${TableName.FolderCommit}.envId`, "=", envId) + .select(selectAllTableCols(TableName.FolderCommit)) + .orderBy("commitId", "desc") + .first(); + return doc; + } catch (error) { + throw new DatabaseError({ error, name: "FindLatestCommit" }); + } + }; + + const findMultipleLatestCommits = async (folderIds: string[], tx?: Knex): Promise => { + try { + const knexInstance = tx || db.replicaNode(); + + // Get the latest commitId for each folderId + const subquery = knexInstance(TableName.FolderCommit) + .whereIn("folderId", folderIds) + .groupBy("folderId") + .select("folderId") + .max("commitId as maxCommitId"); + + // Then fetch the complete rows matching those latest commits + const docs = await knexInstance(TableName.FolderCommit) + // eslint-disable-next-line func-names + .innerJoin(subquery.as("latest"), function () { + this.on(`${TableName.FolderCommit}.folderId`, "=", "latest.folderId").andOn( + `${TableName.FolderCommit}.commitId`, + "=", + "latest.maxCommitId" + ); + }) + .select(selectAllTableCols(TableName.FolderCommit)); + + return docs; + } catch (error) { + throw new DatabaseError({ error, name: "FindMultipleLatestCommits" }); + } + }; + const getNumberOfCommitsSince = async (folderId: string, folderCommitId: string, tx?: Knex): Promise => { try { const referencedCommit = await (tx || db.replicaNode())(TableName.FolderCommit) @@ -62,6 +131,26 @@ export const folderCommitDALFactory = (db: TDbClient) => { } }; + const getEnvNumberOfCommitsSince = async (envId: string, folderCommitId: string, tx?: Knex): Promise => { + try { + const referencedCommit = await (tx || db.replicaNode())(TableName.FolderCommit) + .where({ id: folderCommitId }) + .select("commitId") + .first(); + + if (referencedCommit?.commitId) { + const doc = await (tx || db.replicaNode())(TableName.FolderCommit) + .where(`${TableName.FolderCommit}.envId`, "=", envId) + .where("commitId", ">", referencedCommit.commitId) + .count(); + return Number(doc?.[0].count); + } + return 0; + } catch (error) { + throw new DatabaseError({ error, name: "getNumberOfCommitsSince" }); + } + }; + const findCommitsToRecreate = async ( folderId: string, targetCommitNumber: number, @@ -131,11 +220,77 @@ export const folderCommitDALFactory = (db: TDbClient) => { } }; + const findLatestCommitBetween = async ({ + folderId, + startCommitId, + endCommitId, + tx + }: { + folderId: string; + startCommitId?: string; + endCommitId: string; + tx?: Knex; + }): Promise => { + try { + const doc = await (tx || db.replicaNode())(TableName.FolderCommit) + .where("commitId", "<=", endCommitId) + .where({ folderId }) + .where((qb) => { + if (startCommitId) { + void qb.where("commitId", ">=", startCommitId); + } + }) + .select(selectAllTableCols(TableName.FolderCommit)) + .orderBy("commitId", "desc") + .first(); + return doc; + } catch (error) { + throw new DatabaseError({ error, name: "FindLatestCommitBetween" }); + } + }; + + const findAllCommitsBetween = async ({ + envId, + startCommitId, + endCommitId, + tx + }: { + folderId?: string; + envId?: string; + startCommitId?: string; + endCommitId: string; + tx?: Knex; + }): Promise => { + try { + const docs = await (tx || db.replicaNode())(TableName.FolderCommit) + .where("commitId", "<=", endCommitId) + .where((qb) => { + if (envId) { + void qb.where(`${TableName.FolderCommit}.envId`, "=", envId); + } + if (startCommitId) { + void qb.where("commitId", ">=", startCommitId); + } + }) + .select(selectAllTableCols(TableName.FolderCommit)) + .orderBy("commitId", "desc"); + return docs; + } catch (error) { + throw new DatabaseError({ error, name: "FindLatestCommitBetween" }); + } + }; + return { ...restOfOrm, findByFolderId, findLatestCommit, getNumberOfCommitsSince, - findCommitsToRecreate + findCommitsToRecreate, + findMultipleLatestCommits, + findAllCommitsBetween, + findLatestCommitBetween, + findLatestEnvCommit, + getEnvNumberOfCommitsSince, + findLatestCommitByFolderIds }; }; diff --git a/backend/src/services/folder-commit/folder-commit-queue.ts b/backend/src/services/folder-commit/folder-commit-queue.ts new file mode 100644 index 000000000..06aeaa0c9 --- /dev/null +++ b/backend/src/services/folder-commit/folder-commit-queue.ts @@ -0,0 +1,180 @@ +import { TSecretFolders } from "@app/db/schemas"; +import { getConfig } from "@app/lib/config/env"; +import { logger } from "@app/lib/logger"; +import { QueueJobs, QueueName, TQueueServiceFactory } from "@app/queue"; + +import { TFolderTreeCheckpointDALFactory } from "../folder-tree-checkpoint/folder-tree-checkpoint-dal"; +import { TFolderTreeCheckpointResourcesDALFactory } from "../folder-tree-checkpoint-resources/folder-tree-checkpoint-resources-dal"; +import { TSecretFolderDALFactory } from "../secret-folder/secret-folder-dal"; +import { TFolderCommitDALFactory } from "./folder-commit-dal"; + +type TFolderCommitQueueServiceFactoryDep = { + queueService: TQueueServiceFactory; + folderTreeCheckpointDAL: Pick< + TFolderTreeCheckpointDALFactory, + "create" | "findLatestByEnvId" | "findNearestCheckpoint" + >; + folderTreeCheckpointResourcesDAL: Pick< + TFolderTreeCheckpointResourcesDALFactory, + "insertMany" | "findByTreeCheckpointId" + >; + folderCommitDAL: Pick< + TFolderCommitDALFactory, + "findLatestEnvCommit" | "getEnvNumberOfCommitsSince" | "findMultipleLatestCommits" + >; + folderDAL: Pick; +}; + +export type TFolderCommitQueueServiceFactory = ReturnType; + +export const folderCommitQueueServiceFactory = ({ + queueService, + folderTreeCheckpointDAL, + folderTreeCheckpointResourcesDAL, + folderCommitDAL, + folderDAL +}: TFolderCommitQueueServiceFactoryDep) => { + const appCfg = getConfig(); + + const scheduleTreeCheckpoint = async (envId: string) => { + await queueService.queue( + QueueName.FolderTreeCheckpoint, + QueueJobs.CreateFolderTreeCheckpoint, + { envId }, + { + jobId: envId, + backoff: { + type: "exponential", + delay: 3000 + }, + removeOnFail: { + count: 3 + }, + removeOnComplete: true + } + ); + }; + + const schedulePeriodicTreeCheckpoint = async (envId: string, intervalMs: number) => { + await queueService.queue( + QueueName.FolderTreeCheckpoint, + QueueJobs.CreateFolderTreeCheckpoint, + { envId }, + { + jobId: `periodic-${envId}`, + repeat: { + every: intervalMs + }, + backoff: { + type: "exponential", + delay: 3000 + }, + removeOnFail: false, + removeOnComplete: false + } + ); + }; + + const cancelScheduledTreeCheckpoint = async (envId: string) => { + await queueService.stopJobById(QueueName.FolderTreeCheckpoint, envId); + await queueService.stopRepeatableJobByJobId(QueueName.FolderTreeCheckpoint, `periodic-${envId}`); + }; + + // Sort folders by hierarchy (copied from the source code) + const sortFoldersByHierarchy = (folders: TSecretFolders[]) => { + const childrenMap = new Map(); + const allFolderIds = new Set(); + + folders.forEach((folder) => { + if (folder.id) allFolderIds.add(folder.id); + }); + + folders.forEach((folder) => { + if (folder.parentId) { + const children = childrenMap.get(folder.parentId) || []; + children.push(folder); + childrenMap.set(folder.parentId, children); + } + }); + + const rootFolders = folders.filter((folder) => !folder.parentId || !allFolderIds.has(folder.parentId)); + + const result = []; + let currentLevel = rootFolders; + + while (currentLevel.length > 0) { + result.push(...currentLevel); + + const nextLevel = []; + for (const folder of currentLevel) { + if (folder.id) { + const children = childrenMap.get(folder.id) || []; + nextLevel.push(...children); + } + } + + currentLevel = nextLevel; + } + + return result; + }; + + queueService.start(QueueName.FolderTreeCheckpoint, async (job) => { + try { + if (job.name === QueueJobs.CreateFolderTreeCheckpoint) { + const { envId } = job.data as { envId: string }; + logger.info("Folder tree checkpoint creation started:", envId, job.id); + + const latestTreeCheckpoint = await folderTreeCheckpointDAL.findLatestByEnvId(envId); + + const latestCommit = await folderCommitDAL.findLatestEnvCommit(envId); + if (!latestCommit) { + logger.info(`Latest commit ID not found for envId ${envId}`); + return; + } + const latestCommitId = latestCommit.id; + + if (latestTreeCheckpoint) { + const commitsSinceLastCheckpoint = await folderCommitDAL.getEnvNumberOfCommitsSince( + envId, + latestTreeCheckpoint.folderCommitId + ); + if (commitsSinceLastCheckpoint < Number(appCfg.PIT_TREE_CHECKPOINT_WINDOW)) { + logger.info( + `Commits since last checkpoint ${commitsSinceLastCheckpoint} is less than ${appCfg.PIT_TREE_CHECKPOINT_WINDOW}` + ); + return; + } + } + + const folders = await folderDAL.findByEnvId(envId); + const sortedFolders = sortFoldersByHierarchy(folders); + const filteredFoldersIds = sortedFolders.filter((folder) => !folder.isReserved).map((folder) => folder.id); + + const folderCommits = await folderCommitDAL.findMultipleLatestCommits(filteredFoldersIds); + const folderTreeCheckpoint = await folderTreeCheckpointDAL.create({ + folderCommitId: latestCommitId + }); + + await folderTreeCheckpointResourcesDAL.insertMany( + folderCommits.map((folderCommit) => ({ + folderTreeCheckpointId: folderTreeCheckpoint.id, + folderId: folderCommit.folderId, + folderCommitId: folderCommit.id + })) + ); + + logger.info("Folder tree checkpoint created successfully:", folderTreeCheckpoint.id); + } + } catch (error) { + logger.error(error, "Error creating folder tree checkpoint:"); + throw error; + } + }); + + return { + scheduleTreeCheckpoint, + schedulePeriodicTreeCheckpoint, + cancelScheduledTreeCheckpoint + }; +}; diff --git a/backend/src/services/folder-commit/folder-commit-service.test.ts b/backend/src/services/folder-commit/folder-commit-service.test.ts new file mode 100644 index 000000000..f00e12283 --- /dev/null +++ b/backend/src/services/folder-commit/folder-commit-service.test.ts @@ -0,0 +1,571 @@ +import { Knex } from "knex"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +import { TSecretFolderVersions, TSecretVersionsV2 } from "@app/db/schemas"; +import { BadRequestError, NotFoundError } from "@app/lib/errors"; + +import { ActorType } from "../auth/auth-type"; +import { ChangeType, folderCommitServiceFactory, TFolderCommitServiceFactory } from "./folder-commit-service"; + +// Mock config +vi.mock("@app/lib/config/env", () => ({ + getConfig: () => ({ + PIT_CHECKPOINT_WINDOW: 5, + PIT_TREE_CHECKPOINT_WINDOW: 10 + }) +})); + +// Mock logger +vi.mock("@app/lib/logger", () => ({ + logger: { + info: vi.fn(), + error: vi.fn() + } +})); + +describe("folderCommitServiceFactory", () => { + // Properly type the mock functions + type TransactionCallback = (trx: Knex) => Promise; + + // Mock dependencies + const mockFolderCommitDAL = { + create: vi.fn().mockResolvedValue({}), + findById: vi.fn().mockResolvedValue({}), + findByFolderId: vi.fn().mockResolvedValue([]), + findLatestCommit: vi.fn().mockResolvedValue({}), + transaction: vi.fn().mockImplementation((callback: TransactionCallback) => callback({} as Knex)), + getNumberOfCommitsSince: vi.fn().mockResolvedValue(0), + getEnvNumberOfCommitsSince: vi.fn().mockResolvedValue(0), + findCommitsToRecreate: vi.fn().mockResolvedValue([]), + findMultipleLatestCommits: vi.fn().mockResolvedValue([]), + findLatestCommitBetween: vi.fn().mockResolvedValue({}), + findAllCommitsBetween: vi.fn().mockResolvedValue([]), + findLatestEnvCommit: vi.fn().mockResolvedValue({}), + findLatestCommitByFolderIds: vi.fn().mockResolvedValue({}) + }; + + const mockFolderCommitChangesDAL = { + create: vi.fn().mockResolvedValue({}), + findByCommitId: vi.fn().mockResolvedValue([]), + insertMany: vi.fn().mockResolvedValue([]) + }; + + const mockFolderCheckpointDAL = { + create: vi.fn().mockResolvedValue({}), + findByFolderId: vi.fn().mockResolvedValue([]), + findLatestByFolderId: vi.fn().mockResolvedValue(null), + findNearestCheckpoint: vi.fn().mockResolvedValue({}) + }; + + const mockFolderCheckpointResourcesDAL = { + insertMany: vi.fn().mockResolvedValue([]), + findByCheckpointId: vi.fn().mockResolvedValue([]) + }; + + const mockFolderTreeCheckpointDAL = { + create: vi.fn().mockResolvedValue({}), + findByProjectId: vi.fn().mockResolvedValue([]), + findLatestByProjectId: vi.fn().mockResolvedValue({}), + findNearestCheckpoint: vi.fn().mockResolvedValue({}), + findLatestByEnvId: vi.fn().mockResolvedValue({}) + }; + + const mockFolderTreeCheckpointResourcesDAL = { + insertMany: vi.fn().mockResolvedValue([]), + findByTreeCheckpointId: vi.fn().mockResolvedValue([]) + }; + + const mockUserDAL = { + findById: vi.fn().mockResolvedValue({}) + }; + + const mockIdentityDAL = { + findById: vi.fn().mockResolvedValue({}) + }; + + const mockFolderDAL = { + findByParentId: vi.fn().mockResolvedValue([]), + findByProjectId: vi.fn().mockResolvedValue([]), + deleteById: vi.fn().mockResolvedValue({}), + create: vi.fn().mockResolvedValue({}), + updateById: vi.fn().mockResolvedValue({}), + update: vi.fn().mockResolvedValue({}), + find: vi.fn().mockResolvedValue([]), + findById: vi.fn().mockResolvedValue({}), + findByEnvId: vi.fn().mockResolvedValue([]), + findFoldersByRootAndIds: vi.fn().mockResolvedValue([]) + }; + + const mockFolderVersionDAL = { + findLatestFolderVersions: vi.fn().mockResolvedValue({}), + findById: vi.fn().mockResolvedValue({}), + deleteById: vi.fn().mockResolvedValue({}), + create: vi.fn().mockResolvedValue({}), + updateById: vi.fn().mockResolvedValue({}), + find: vi.fn().mockResolvedValue([]), + findByIdsWithLatestVersion: vi.fn().mockResolvedValue({}) + }; + + const mockSecretVersionV2BridgeDAL = { + findLatestVersionByFolderId: vi.fn().mockResolvedValue([]), + findById: vi.fn().mockResolvedValue({}), + deleteById: vi.fn().mockResolvedValue({}), + create: vi.fn().mockResolvedValue({}), + updateById: vi.fn().mockResolvedValue({}), + find: vi.fn().mockResolvedValue([]), + findByIdsWithLatestVersion: vi.fn().mockResolvedValue({}) + }; + + const mockSecretV2BridgeDAL = { + deleteById: vi.fn().mockResolvedValue({}), + create: vi.fn().mockResolvedValue({}), + updateById: vi.fn().mockResolvedValue({}), + update: vi.fn().mockResolvedValue({}), + insertMany: vi.fn().mockResolvedValue([]), + invalidateSecretCacheByProjectId: vi.fn().mockResolvedValue({}) + }; + + const mockProjectDAL = { + findById: vi.fn().mockResolvedValue({}) + }; + + const mockFolderCommitQueueService = { + scheduleTreeCheckpoint: vi.fn().mockResolvedValue({}) + }; + + let folderCommitService: TFolderCommitServiceFactory; + + beforeEach(() => { + vi.clearAllMocks(); + + folderCommitService = folderCommitServiceFactory({ + folderCommitDAL: mockFolderCommitDAL, + folderCommitChangesDAL: mockFolderCommitChangesDAL, + folderCheckpointDAL: mockFolderCheckpointDAL, + folderCheckpointResourcesDAL: mockFolderCheckpointResourcesDAL, + folderTreeCheckpointDAL: mockFolderTreeCheckpointDAL, + folderTreeCheckpointResourcesDAL: mockFolderTreeCheckpointResourcesDAL, + userDAL: mockUserDAL, + identityDAL: mockIdentityDAL, + folderDAL: mockFolderDAL, + folderVersionDAL: mockFolderVersionDAL, + secretVersionV2BridgeDAL: mockSecretVersionV2BridgeDAL, + projectDAL: mockProjectDAL, + secretV2BridgeDAL: mockSecretV2BridgeDAL, + folderCommitQueueService: mockFolderCommitQueueService + }); + }); + + afterEach(() => { + vi.resetAllMocks(); + }); + + describe("createCommit", () => { + it("should successfully create a commit with user actor", async () => { + // Arrange + const userData = { id: "user-id", username: "testuser" }; + const folderData = { id: "folder-id", envId: "env-id" }; + const commitData = { id: "commit-id", folderId: "folder-id" }; + + mockUserDAL.findById.mockResolvedValue(userData); + mockFolderDAL.findById.mockResolvedValue(folderData); + mockFolderCommitDAL.create.mockResolvedValue(commitData); + mockFolderCheckpointDAL.findLatestByFolderId.mockResolvedValue(null); + mockFolderCommitDAL.findLatestCommit.mockResolvedValue({ id: "latest-commit-id" }); + mockFolderDAL.findByParentId.mockResolvedValue([]); + mockSecretVersionV2BridgeDAL.findLatestVersionByFolderId.mockResolvedValue([]); + + const data = { + actor: { + type: ActorType.USER, + metadata: { id: userData.id } + }, + message: "Test commit", + folderId: folderData.id, + changes: [ + { + type: "add", + secretVersionId: "secret-version-1" + } + ] + }; + + // Act + const result = await folderCommitService.createCommit(data); + + // Assert + expect(mockUserDAL.findById).toHaveBeenCalledWith(userData.id, undefined); + expect(mockFolderDAL.findById).toHaveBeenCalledWith(folderData.id, undefined); + expect(mockFolderCommitDAL.create).toHaveBeenCalledWith( + expect.objectContaining({ + actorType: ActorType.USER, + // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment + actorMetadata: expect.objectContaining({ name: userData.username }), + message: data.message, + folderId: data.folderId, + envId: folderData.envId + }), + undefined + ); + expect(mockFolderCommitChangesDAL.insertMany).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ + folderCommitId: commitData.id, + changeType: data.changes[0].type, + secretVersionId: data.changes[0].secretVersionId + }) + ]), + undefined + ); + expect(mockFolderCommitQueueService.scheduleTreeCheckpoint).toHaveBeenCalledWith(folderData.envId); + expect(result).toEqual(commitData); + }); + + it("should successfully create a commit with identity actor", async () => { + // Arrange + const identityData = { id: "identity-id", name: "testidentity" }; + const folderData = { id: "folder-id", envId: "env-id" }; + const commitData = { id: "commit-id", folderId: "folder-id" }; + + mockIdentityDAL.findById.mockResolvedValue(identityData); + mockFolderDAL.findById.mockResolvedValue(folderData); + mockFolderCommitDAL.create.mockResolvedValue(commitData); + mockFolderCheckpointDAL.findLatestByFolderId.mockResolvedValue(null); + mockFolderCommitDAL.findLatestCommit.mockResolvedValue({ id: "latest-commit-id" }); + mockFolderDAL.findByParentId.mockResolvedValue([]); + mockSecretVersionV2BridgeDAL.findLatestVersionByFolderId.mockResolvedValue([]); + + const data = { + actor: { + type: ActorType.IDENTITY, + metadata: { id: identityData.id } + }, + message: "Test commit", + folderId: folderData.id, + changes: [ + { + type: "add", + folderVersionId: "folder-version-1" + } + ] + }; + + // Act + const result = await folderCommitService.createCommit(data); + + // Assert + expect(mockIdentityDAL.findById).toHaveBeenCalledWith(identityData.id, undefined); + expect(mockFolderDAL.findById).toHaveBeenCalledWith(folderData.id, undefined); + expect(mockFolderCommitDAL.create).toHaveBeenCalledWith( + expect.objectContaining({ + actorType: ActorType.IDENTITY, + // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment + actorMetadata: expect.objectContaining({ name: identityData.name }), + message: data.message, + folderId: data.folderId, + envId: folderData.envId + }), + undefined + ); + expect(mockFolderCommitChangesDAL.insertMany).toHaveBeenCalledWith( + expect.arrayContaining([ + expect.objectContaining({ + folderCommitId: commitData.id, + changeType: data.changes[0].type, + folderVersionId: data.changes[0].folderVersionId + }) + ]), + undefined + ); + expect(mockFolderCommitQueueService.scheduleTreeCheckpoint).toHaveBeenCalledWith(folderData.envId); + expect(result).toEqual(commitData); + }); + + it("should throw NotFoundError when folder does not exist", async () => { + // Arrange + mockFolderDAL.findById.mockResolvedValue(null); + + const data = { + actor: { + type: ActorType.PLATFORM + }, + message: "Test commit", + folderId: "non-existent-folder", + changes: [] + }; + + // Act & Assert + await expect(folderCommitService.createCommit(data)).rejects.toThrow(NotFoundError); + expect(mockFolderDAL.findById).toHaveBeenCalledWith("non-existent-folder", undefined); + }); + }); + + describe("addCommitChange", () => { + it("should successfully add a change to an existing commit", async () => { + // Arrange + const commitData = { id: "commit-id", folderId: "folder-id" }; + const changeData = { id: "change-id", folderCommitId: "commit-id" }; + + mockFolderCommitDAL.findById.mockResolvedValue(commitData); + mockFolderCommitChangesDAL.create.mockResolvedValue(changeData); + + const data = { + folderCommitId: commitData.id, + changeType: "add", + secretVersionId: "secret-version-1" + }; + + // Act + const result = await folderCommitService.addCommitChange(data); + + // Assert + expect(mockFolderCommitDAL.findById).toHaveBeenCalledWith(commitData.id, undefined); + expect(mockFolderCommitChangesDAL.create).toHaveBeenCalledWith(data, undefined); + expect(result).toEqual(changeData); + }); + + it("should throw BadRequestError when neither secretVersionId nor folderVersionId is provided", async () => { + // Arrange + const data = { + folderCommitId: "commit-id", + changeType: "add" + }; + + // Act & Assert + await expect(folderCommitService.addCommitChange(data)).rejects.toThrow(BadRequestError); + }); + + it("should throw NotFoundError when commit does not exist", async () => { + // Arrange + mockFolderCommitDAL.findById.mockResolvedValue(null); + + const data = { + folderCommitId: "non-existent-commit", + changeType: "add", + secretVersionId: "secret-version-1" + }; + + // Act & Assert + await expect(folderCommitService.addCommitChange(data)).rejects.toThrow(NotFoundError); + expect(mockFolderCommitDAL.findById).toHaveBeenCalledWith("non-existent-commit", undefined); + }); + }); + + // Note: reconstructFolderState is an internal function not exposed in the public API + // We'll test it indirectly through compareFolderStates + + describe("compareFolderStates", () => { + it("should mark all resources as creates when currentCommitId is not provided", async () => { + // Arrange + const targetCommitId = "target-commit-id"; + const targetCommit = { id: targetCommitId, commitId: 1, folderId: "folder-id" }; + + mockFolderCommitDAL.findById.mockResolvedValue(targetCommit); + // Mock how compareFolderStates would process the results internally + mockFolderCheckpointDAL.findNearestCheckpoint.mockResolvedValue({ id: "checkpoint-id", commitId: "hash-0" }); + mockFolderCheckpointResourcesDAL.findByCheckpointId.mockResolvedValue([ + { secretVersionId: "secret-version-1", referencedSecretId: "secret-1" }, + { folderVersionId: "folder-version-1", referencedFolderId: "folder-1" } + ]); + mockFolderCommitDAL.findCommitsToRecreate.mockResolvedValue([]); + + // Act + const result = await folderCommitService.compareFolderStates({ + targetCommitId + }); + + // Assert + expect(mockFolderCommitDAL.findById).toHaveBeenCalledWith(targetCommitId, undefined); + + // Verify we get resources marked as create + expect(result).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + changeType: "create", + commitId: targetCommit.commitId + }) + ]) + ); + }); + }); + + describe("createFolderCheckpoint", () => { + it("should successfully create a checkpoint when force is true", async () => { + // Arrange + const folderCommitId = "commit-id"; + const folderId = "folder-id"; + const checkpointData = { id: "checkpoint-id", folderCommitId }; + + mockFolderDAL.findByParentId.mockResolvedValue([{ id: "subfolder-id" }]); + mockFolderVersionDAL.findLatestFolderVersions.mockResolvedValue({ "subfolder-id": { id: "folder-version-1" } }); + mockSecretVersionV2BridgeDAL.findLatestVersionByFolderId.mockResolvedValue([{ id: "secret-version-1" }]); + mockFolderCheckpointDAL.create.mockResolvedValue(checkpointData); + + // Act + const result = await folderCommitService.createFolderCheckpoint({ + folderId, + folderCommitId, + force: true + }); + + // Assert + expect(mockFolderCheckpointDAL.create).toHaveBeenCalledWith({ folderCommitId }, undefined); + expect(mockFolderCheckpointResourcesDAL.insertMany).toHaveBeenCalled(); + expect(result).toBe(folderCommitId); + }); + }); + + describe("deepRollbackFolder", () => { + it("should throw NotFoundError when commit doesn't exist", async () => { + // Arrange + const targetCommitId = "non-existent-commit"; + const envId = "env-id"; + const actorId = "user-id"; + const actorType = ActorType.USER; + const projectId = "project-id"; + + mockFolderCommitDAL.findById.mockResolvedValue(null); + + // Act & Assert + await expect( + folderCommitService.deepRollbackFolder(targetCommitId, envId, actorId, actorType, projectId) + ).rejects.toThrow(NotFoundError); + }); + }); + + describe("createFolderTreeCheckpoint", () => { + it("should create a tree checkpoint when checkpoint window is exceeded", async () => { + // Arrange + const envId = "env-id"; + const folderCommitId = "commit-id"; + const latestCommit = { id: folderCommitId }; + const latestTreeCheckpoint = { id: "tree-checkpoint-id", folderCommitId: "old-commit-id" }; + const folders = [ + { id: "folder-1", isReserved: false }, + { id: "folder-2", isReserved: false }, + { id: "folder-3", isReserved: true } // Reserved folders should be filtered out + ]; + const folderCommits = [ + { folderId: "folder-1", id: "commit-1" }, + { folderId: "folder-2", id: "commit-2" } + ]; + const treeCheckpoint = { id: "new-tree-checkpoint-id" }; + + mockFolderCommitDAL.findLatestEnvCommit.mockResolvedValue(latestCommit); + mockFolderTreeCheckpointDAL.findLatestByEnvId.mockResolvedValue(latestTreeCheckpoint); + mockFolderCommitDAL.getEnvNumberOfCommitsSince.mockResolvedValue(15); // More than PIT_TREE_CHECKPOINT_WINDOW (10) + mockFolderDAL.findByEnvId.mockResolvedValue(folders); + mockFolderCommitDAL.findMultipleLatestCommits.mockResolvedValue(folderCommits); + mockFolderTreeCheckpointDAL.create.mockResolvedValue(treeCheckpoint); + + // Act + await folderCommitService.createFolderTreeCheckpoint(envId); + + // Assert + expect(mockFolderCommitDAL.findLatestEnvCommit).toHaveBeenCalledWith(envId, undefined); + expect(mockFolderTreeCheckpointDAL.create).toHaveBeenCalledWith({ folderCommitId }, undefined); + }); + }); + + describe("applyFolderStateDifferences", () => { + it("should process changes correctly", async () => { + // Arrange + const folderId = "folder-id"; + const projectId = "project-id"; + const actorId = "user-id"; + const actorType = ActorType.USER; + + const differences = [ + { type: "secret", id: "secret-1", versionId: "v1", changeType: ChangeType.CREATE, commitId: 1 }, + { type: "folder", id: "folder-1", versionId: "v2", changeType: ChangeType.UPDATE, commitId: 1 } + ]; + + const secretVersions = { + "secret-1": { + id: "secret-version-1", + createdAt: new Date(), + updatedAt: new Date(), + type: "shared", + folderId: "folder-1", + secretId: "secret-1", + version: 1, + key: "SECRET_KEY", + encryptedValue: Buffer.from("encrypted"), + encryptedComment: Buffer.from("comment"), + skipMultilineEncoding: false, + userId: "user-1", + envId: "env-1", + metadata: {} + } as TSecretVersionsV2 + }; + + const folderVersions = { + "folder-1": { + folderId: "folder-1", + version: 1, + name: "Test Folder", + envId: "env-1" + } as TSecretFolderVersions + }; + + // Mock folder lookup for the folder being processed + mockFolderDAL.findById.mockImplementation((id) => { + if (id === folderId) { + return Promise.resolve({ id: folderId, envId: "env-1" }); + } + return Promise.resolve(null); + }); + + // Mock latest commit lookup + mockFolderCommitDAL.findLatestCommit.mockImplementation((id) => { + if (id === folderId) { + return Promise.resolve({ id: "latest-commit-id", folderId }); + } + return Promise.resolve(null); + }); + + // Make sure findByParentId returns an array, not undefined + mockFolderDAL.findByParentId.mockResolvedValue([]); + + // Make sure other required functions return appropriate values + mockFolderCheckpointDAL.findLatestByFolderId.mockResolvedValue(null); + mockSecretVersionV2BridgeDAL.findLatestVersionByFolderId.mockResolvedValue([]); + + // These mocks need to return objects with an id field + mockSecretVersionV2BridgeDAL.findByIdsWithLatestVersion.mockResolvedValue(secretVersions); + mockFolderVersionDAL.findByIdsWithLatestVersion.mockResolvedValue(folderVersions); + mockSecretV2BridgeDAL.insertMany.mockResolvedValue([{ id: "new-secret-1" }]); + mockSecretVersionV2BridgeDAL.create.mockResolvedValue({ id: "new-secret-version-1" }); + mockFolderDAL.updateById.mockResolvedValue({ id: "updated-folder-1" }); + mockFolderVersionDAL.create.mockResolvedValue({ id: "new-folder-version-1" }); + mockFolderCommitDAL.create.mockResolvedValue({ id: "new-commit-id" }); + + // Mock transaction + mockFolderCommitDAL.transaction.mockImplementation((callback: TransactionCallback) => callback({} as Knex)); + + // Act + const result = await folderCommitService.applyFolderStateDifferences({ + differences, + actorInfo: { + actorType, + actorId, + message: "Applying changes" + }, + folderId, + projectId, + reconstructNewFolders: false + }); + + // Assert + expect(mockFolderCommitDAL.create).toHaveBeenCalled(); + expect(mockSecretV2BridgeDAL.invalidateSecretCacheByProjectId).toHaveBeenCalledWith(projectId); + + // Check that we got the right counts + expect(result).toEqual({ + secretChangesCount: 1, + folderChangesCount: 1, + totalChanges: 2 + }); + }); + }); +}); diff --git a/backend/src/services/folder-commit/folder-commit-service.ts b/backend/src/services/folder-commit/folder-commit-service.ts index e40847d06..1fc4b456d 100644 --- a/backend/src/services/folder-commit/folder-commit-service.ts +++ b/backend/src/services/folder-commit/folder-commit-service.ts @@ -1,23 +1,87 @@ /* eslint-disable no-await-in-loop */ import { Knex } from "knex"; -import { TSecretFolders } from "@app/db/schemas"; +import { TSecretFolders, TSecretFolderVersions, TSecretVersionsV2 } from "@app/db/schemas"; import { getConfig } from "@app/lib/config/env"; import { BadRequestError, DatabaseError, NotFoundError } from "@app/lib/errors"; +import { logger } from "@app/lib/logger"; import { ActorType } from "../auth/auth-type"; import { TFolderCheckpointDALFactory } from "../folder-checkpoint/folder-checkpoint-dal"; import { TFolderCheckpointResourcesDALFactory } from "../folder-checkpoint-resources/folder-checkpoint-resources-dal"; import { TFolderCommitChangesDALFactory } from "../folder-commit-changes/folder-commit-changes-dal"; import { TFolderTreeCheckpointDALFactory } from "../folder-tree-checkpoint/folder-tree-checkpoint-dal"; +import { TFolderTreeCheckpointResourcesDALFactory } from "../folder-tree-checkpoint-resources/folder-tree-checkpoint-resources-dal"; import { TIdentityDALFactory } from "../identity/identity-dal"; import { TProjectDALFactory } from "../project/project-dal"; import { TSecretFolderDALFactory } from "../secret-folder/secret-folder-dal"; import { TSecretFolderVersionDALFactory } from "../secret-folder/secret-folder-version-dal"; -import { TSecretV2BridgeDALFactory } from "../secret-v2-bridge/secret-v2-bridge-dal"; +import * as secretV2BridgeDal from "../secret-v2-bridge/secret-v2-bridge-dal"; import { TSecretVersionV2DALFactory } from "../secret-v2-bridge/secret-version-dal"; import { TUserDALFactory } from "../user/user-dal"; import { TFolderCommitDALFactory } from "./folder-commit-dal"; +import { TFolderCommitQueueServiceFactory } from "./folder-commit-queue"; + +// Define enums for better type safety +export enum ChangeType { + ADD = "add", + DELETE = "delete", + UPDATE = "update", + CREATE = "create" +} + +enum ResourceType { + SECRET = "secret", + FOLDER = "folder" +} + +// Improved types for DTO objects +type TCreateCommitDTO = { + actor: { + type: string; + metadata?: { + name?: string; + id?: string; + }; + }; + message?: string; + folderId: string; + changes: { + type: string; + secretVersionId?: string; + folderVersionId?: string; + }[]; +}; + +type TCommitChangeDTO = { + folderCommitId: string; + changeType: string; + secretVersionId?: string; + folderVersionId?: string; +}; + +export type ResourceChange = { + type: string; + id: string; + versionId: string; + oldVersionId?: string; + changeType: ChangeType; + commitId: number; + createdAt?: Date; + parentId?: string; +}; + +type ActorInfo = { + actorType: string; + actorId?: string; + message?: string; +}; + +type StateChangeResult = { + secretChangesCount: number; + folderChangesCount: number; + totalChanges: number; +}; type TFolderCommitServiceFactoryDep = { folderCommitDAL: Pick< @@ -28,7 +92,13 @@ type TFolderCommitServiceFactoryDep = { | "findLatestCommit" | "transaction" | "getNumberOfCommitsSince" + | "getEnvNumberOfCommitsSince" | "findCommitsToRecreate" + | "findMultipleLatestCommits" + | "findLatestCommitBetween" + | "findAllCommitsBetween" + | "findLatestEnvCommit" + | "findLatestCommitByFolderIds" >; folderCommitChangesDAL: Pick; folderCheckpointDAL: Pick< @@ -36,12 +106,28 @@ type TFolderCommitServiceFactoryDep = { "create" | "findByFolderId" | "findLatestByFolderId" | "findNearestCheckpoint" >; folderCheckpointResourcesDAL: Pick; - folderTreeCheckpointDAL: Pick; + folderTreeCheckpointDAL: Pick< + TFolderTreeCheckpointDALFactory, + "create" | "findNearestCheckpoint" | "findLatestByEnvId" + >; + folderTreeCheckpointResourcesDAL: Pick< + TFolderTreeCheckpointResourcesDALFactory, + "insertMany" | "findByTreeCheckpointId" + >; userDAL: Pick; identityDAL: Pick; folderDAL: Pick< TSecretFolderDALFactory, - "findByParentId" | "findByProjectId" | "deleteById" | "create" | "updateById" | "update" + | "findByParentId" + | "findByProjectId" + | "deleteById" + | "create" + | "updateById" + | "update" + | "find" + | "findById" + | "findByEnvId" + | "findFoldersByRootAndIds" >; folderVersionDAL: Pick< TSecretFolderVersionDALFactory, @@ -64,34 +150,11 @@ type TFolderCommitServiceFactoryDep = { | "findByIdsWithLatestVersion" >; secretV2BridgeDAL: Pick< - TSecretV2BridgeDALFactory, + secretV2BridgeDal.TSecretV2BridgeDALFactory, "deleteById" | "create" | "updateById" | "update" | "insertMany" | "invalidateSecretCacheByProjectId" >; projectDAL: Pick; -}; - -export type TCreateCommitDTO = { - actor: { - type: string; - metadata?: { - name?: string; - id?: string; - }; - }; - message?: string; - folderId: string; - changes: { - type: string; - secretVersionId?: string; - folderVersionId?: string; - }[]; -}; - -export type TCommitChangeDTO = { - folderCommitId: string; - changeType: string; - secretVersionId?: string; - folderVersionId?: string; + folderCommitQueueService?: Pick; }; export const folderCommitServiceFactory = ({ @@ -106,25 +169,36 @@ export const folderCommitServiceFactory = ({ folderVersionDAL, secretVersionV2BridgeDAL, projectDAL, - secretV2BridgeDAL + secretV2BridgeDAL, + folderTreeCheckpointResourcesDAL, + folderCommitQueueService }: TFolderCommitServiceFactoryDep) => { const appCfg = getConfig(); + /** + * Fetches all resources within a folder + */ const getFolderResources = async (folderId: string, tx?: Knex) => { const resources = []; const subFolders = await folderDAL.findByParentId(folderId, tx); + if (subFolders.length > 0) { const subFolderIds = subFolders.map((folder) => folder.id); const folderVersions = await folderVersionDAL.findLatestFolderVersions(subFolderIds, tx); resources.push(...Object.values(folderVersions).map((folderVersion) => ({ folderVersionId: folderVersion.id }))); } + const secretVersions = await secretVersionV2BridgeDAL.findLatestVersionByFolderId(folderId, tx); if (secretVersions.length > 0) { resources.push(...secretVersions.map((secretVersion) => ({ secretVersionId: secretVersion.id }))); } + return resources; }; + /** + * Creates a checkpoint for a folder if necessary + */ const createFolderCheckpoint = async ({ folderId, folderCommitId, @@ -138,23 +212,26 @@ export const folderCommitServiceFactory = ({ }) => { let latestCommitId = folderCommitId; const latestCheckpoint = await folderCheckpointDAL.findLatestByFolderId(folderId, tx); + if (!latestCommitId) { - latestCommitId = (await folderCommitDAL.findLatestCommit(folderId, tx))?.id; - } - if (!latestCommitId) { - throw new BadRequestError({ message: "Latest commit ID not found" }); - return; + const latestCommit = await folderCommitDAL.findLatestCommit(folderId, tx); + if (!latestCommit) { + throw new BadRequestError({ message: "Latest commit ID not found" }); + } + latestCommitId = latestCommit.id; } + if (!force && latestCheckpoint) { const commitsSinceLastCheckpoint = await folderCommitDAL.getNumberOfCommitsSince( folderId, latestCheckpoint.folderCommitId, tx ); - if (commitsSinceLastCheckpoint < Number(appCfg.CHECKPOINT_WINDOW)) { + if (commitsSinceLastCheckpoint < Number(appCfg.PIT_CHECKPOINT_WINDOW)) { return; } } + const checkpointResources = await getFolderResources(folderId, tx); if (checkpointResources.length > 0) { @@ -169,8 +246,12 @@ export const folderCommitServiceFactory = ({ tx ); } + return latestCommitId; }; + /** + * Reconstructs the state of a folder at a specific commit + */ const reconstructFolderState = async ( folderCommitId: string, tx?: Knex @@ -249,8 +330,37 @@ export const folderCommitServiceFactory = ({ return Object.values(folderState); }; - const compareFolderStates = async (currentCommitId: string, targetCommitId: string, tx?: Knex) => { - // Reconstruct state for both commits + /** + * Compares folder states between two commits and returns the differences + */ + const compareFolderStates = async ({ + currentCommitId, + targetCommitId, + tx + }: { + currentCommitId?: string; + targetCommitId: string; + tx?: Knex; + }) => { + const targetCommit = await folderCommitDAL.findById(targetCommitId, tx); + if (!targetCommit) { + throw new NotFoundError({ message: `Commit with ID ${targetCommitId} not found` }); + } + + // If currentCommitId is not provided, mark all resources in target as creates + if (!currentCommitId) { + const targetState = await reconstructFolderState(targetCommitId, tx); + + return targetState.map((resource) => ({ + type: resource.type, + id: resource.id, + versionId: resource.versionId, + changeType: "create", + commitId: targetCommit.commitId + })) as ResourceChange[]; + } + + // Original logic for when currentCommitId is provided const currentState = await reconstructFolderState(currentCommitId, tx); const targetState = await reconstructFolderState(targetCommitId, tx); @@ -258,52 +368,45 @@ export const folderCommitServiceFactory = ({ const currentMap: Record = {}; const targetMap: Record = {}; - // Build lookup map for current state + // Build lookup maps currentState.forEach((resource) => { const key = `${resource.type}-${resource.id}`; currentMap[key] = resource; }); - // Build lookup map for target state targetState.forEach((resource) => { const key = `${resource.type}-${resource.id}`; targetMap[key] = resource; }); // Track differences - const differences: { - type: string; - id: string; - versionId: string; - changeType: "create" | "update" | "delete"; - }[] = []; + const differences: ResourceChange[] = []; - // Find deletes and updates (resources in current but not in target, or with different versions) + // Find deletes and updates Object.keys(currentMap).forEach((key) => { const currentResource = currentMap[key]; const targetResource = targetMap[key]; if (!targetResource) { - // Resource exists in current but not in target - it's a delete differences.push({ type: currentResource.type, id: currentResource.id, versionId: currentResource.versionId, - changeType: "delete" + changeType: ChangeType.DELETE, + commitId: targetCommit.commitId }); } else if (currentResource.versionId !== targetResource.versionId) { - // Resource exists in both but with different versions - it's an update differences.push({ type: targetResource.type, id: targetResource.id, versionId: targetResource.versionId, - changeType: "update" + changeType: ChangeType.UPDATE, + commitId: targetCommit.commitId }); } - // If versions are the same, it's unchanged - exclude from result }); - // Find creates (resources in target but not in current) + // Find creates Object.keys(targetMap).forEach((key) => { if (!currentMap[key]) { const targetResource = targetMap[key]; @@ -311,7 +414,9 @@ export const folderCommitServiceFactory = ({ type: targetResource.type, id: targetResource.id, versionId: targetResource.versionId, - changeType: "create" + changeType: ChangeType.CREATE, + commitId: targetCommit.commitId, + createdAt: targetCommit.createdAt }); } }); @@ -319,7 +424,9 @@ export const folderCommitServiceFactory = ({ return differences; }; - // Add a change to an existing commit + /** + * Adds a change to an existing commit + */ const addCommitChange = async (data: TCommitChangeDTO, tx?: Knex) => { try { if (!data.secretVersionId && !data.folderVersionId) { @@ -340,26 +447,39 @@ export const folderCommitServiceFactory = ({ } }; + /** + * Creates a new commit with the provided changes + */ const createCommit = async (data: TCreateCommitDTO, tx?: Knex) => { - const metadata = data.actor.metadata || {}; try { + const metadata = { ...data.actor.metadata } || {}; + if (data.actor.type === ActorType.USER && data.actor.metadata?.id) { const user = await userDAL.findById(data.actor.metadata?.id, tx); metadata.name = user?.username; } + if (data.actor.type === ActorType.IDENTITY && data.actor.metadata?.id) { const identity = await identityDAL.findById(data.actor.metadata?.id, tx); metadata.name = identity?.name; } + + const folder = await folderDAL.findById(data.folderId, tx); + if (!folder) { + throw new NotFoundError({ message: `Folder with ID ${data.folderId} not found` }); + } + const newCommit = await folderCommitDAL.create( { actorMetadata: metadata, actorType: data.actor.type, message: data.message, - folderId: data.folderId + folderId: data.folderId, + envId: folder.envId }, tx ); + await folderCommitChangesDAL.insertMany( data.changes.map((change) => ({ folderCommitId: newCommit.id, @@ -371,167 +491,184 @@ export const folderCommitServiceFactory = ({ ); await createFolderCheckpoint({ folderId: data.folderId, tx }); + if (folderCommitQueueService) { + await folderCommitQueueService.scheduleTreeCheckpoint(folder.envId); + } return newCommit; } catch (error) { + if (error instanceof NotFoundError || error instanceof BadRequestError) { + throw error; + } throw new DatabaseError({ error, name: "CreateCommit" }); } }; - const applyFolderStateDifferences = async ( - differences: Array<{ - type: string; - id: string; - versionId: string; - oldVersionId?: string; - changeType: "create" | "update" | "delete"; - }>, - actorInfo: { - actorType: string; - actorId?: string; - message?: string; - }, + /** + * Process secret changes when applying folder state differences + */ + const processSecretChanges = async ( + changes: ResourceChange[], + secretVersions: Record, + actorInfo: ActorInfo, folderId: string, - projectId: string + tx?: Knex ) => { - let result = {}; - await folderCommitDAL.transaction(async (tx) => { - // Group differences by type for more efficient processing - const secretChanges = differences.filter((diff) => diff.type === "secret"); - const folderChanges = differences.filter((diff) => diff.type === "folder"); + const commitChanges = []; - const secretVersions = await secretVersionV2BridgeDAL.findByIdsWithLatestVersion( - folderId, - secretChanges.map((diff) => diff.id), - secretChanges.map((diff) => diff.versionId) - ); - const folderVersions = await folderVersionDAL.findByIdsWithLatestVersion( - folderChanges.map((diff) => diff.id), - folderChanges.map((diff) => diff.versionId) - ); + for (const change of changes) { + const secretVersion = secretVersions[change.id]; + // eslint-disable-next-line no-continue + if (!secretVersion) continue; - // Track all changes for commit recording - const commitChanges = []; + switch (change.changeType) { + case "create": + { + const newSecret = [ + { + id: change.id, + skipMultilineEncoding: secretVersion.skipMultilineEncoding, + version: secretVersion.version + 1, + type: secretVersion.type, + key: secretVersion.key, + reminderNote: secretVersion.reminderNote, + reminderRepeatDays: secretVersion.reminderRepeatDays, + encryptedValue: secretVersion.encryptedValue, + encryptedComment: secretVersion.encryptedComment, + userId: secretVersion.userId, + metadata: secretVersion.metadata, + folderId + } + ]; + await secretV2BridgeDAL.insertMany(newSecret, tx); - // Process secret changes - for (const change of secretChanges) { - const secretVersion = secretVersions[change.id]; - switch (change.changeType) { - case "create": - if (secretVersion) { - const newSecret = [ - { - id: change.id, - skipMultilineEncoding: secretVersion.skipMultilineEncoding, - version: secretVersion.version + 1, - type: secretVersion.type, - key: secretVersion.key, - reminderNote: secretVersion.reminderNote, - reminderRepeatDays: secretVersion.reminderRepeatDays, - encryptedValue: secretVersion.encryptedValue, - encryptedComment: secretVersion.encryptedComment, - userId: secretVersion.userId, - metadata: secretVersion.metadata, - folderId - } - ]; - await secretV2BridgeDAL.insertMany(newSecret, tx); - - const newVersion = await secretVersionV2BridgeDAL.create( - { - folderId, - secretId: secretVersion.secretId, - version: secretVersion.version + 1, - encryptedValue: secretVersion.encryptedValue, - key: secretVersion.key, - encryptedComment: secretVersion.encryptedComment, - skipMultilineEncoding: secretVersion.skipMultilineEncoding, - reminderNote: secretVersion.reminderNote, - reminderRepeatDays: secretVersion.reminderRepeatDays, - userId: secretVersion.userId, - metadata: secretVersion.metadata, - actorType: actorInfo.actorType, - envId: secretVersion.envId, - ...(actorInfo.actorType === ActorType.IDENTITY && { identityActorId: actorInfo.actorId }), - ...(actorInfo.actorType === ActorType.USER && { userActorId: actorInfo.actorId }) - }, - tx - ); - - commitChanges.push({ - type: "add", - secretVersionId: newVersion.id - }); - } - break; - - case "update": - // Update secret to specific version - if (secretVersion) { - await secretV2BridgeDAL.updateById( - change.id, - { - skipMultilineEncoding: secretVersion?.skipMultilineEncoding, - version: secretVersion?.version, - type: secretVersion?.type, - key: secretVersion?.key, - reminderNote: secretVersion?.reminderNote, - reminderRepeatDays: secretVersion?.reminderRepeatDays, - encryptedValue: secretVersion?.encryptedValue, - encryptedComment: secretVersion?.encryptedComment, - userId: secretVersion?.userId, - metadata: secretVersion?.metadata - }, - tx - ); - - const newVersion = await secretVersionV2BridgeDAL.create( - { - version: secretVersion.version + 1, - encryptedValue: secretVersion.encryptedValue, - key: secretVersion.key, - encryptedComment: secretVersion.encryptedComment, - skipMultilineEncoding: secretVersion.skipMultilineEncoding, - reminderNote: secretVersion.reminderNote, - reminderRepeatDays: secretVersion.reminderRepeatDays, - userId: secretVersion.userId, - metadata: secretVersion.metadata, - actorType: actorInfo.actorType, - envId: secretVersion.envId, - folderId, - secretId: secretVersion.secretId, - ...(actorInfo.actorType === ActorType.IDENTITY && { identityActorId: actorInfo.actorId }), - ...(actorInfo.actorType === ActorType.USER && { userActorId: actorInfo.actorId }) - }, - tx - ); - - commitChanges.push({ - type: "add", - secretVersionId: newVersion.id - }); - } - break; - - case "delete": - await secretV2BridgeDAL.deleteById(change.id, tx); + const newVersion = await secretVersionV2BridgeDAL.create( + { + folderId, + secretId: secretVersion.secretId, + version: secretVersion.version + 1, + encryptedValue: secretVersion.encryptedValue, + key: secretVersion.key, + encryptedComment: secretVersion.encryptedComment, + skipMultilineEncoding: secretVersion.skipMultilineEncoding, + reminderNote: secretVersion.reminderNote, + reminderRepeatDays: secretVersion.reminderRepeatDays, + userId: secretVersion.userId, + metadata: secretVersion.metadata, + actorType: actorInfo.actorType, + envId: secretVersion.envId, + ...(actorInfo.actorType === ActorType.IDENTITY && { identityActorId: actorInfo.actorId }), + ...(actorInfo.actorType === ActorType.USER && { userActorId: actorInfo.actorId }) + }, + tx + ); commitChanges.push({ - type: "delete", - secretVersionId: change.versionId + type: ChangeType.ADD, + secretVersionId: newVersion.id }); - break; + } + break; - default: - throw new BadRequestError({ message: `Unknown change type: ${change.changeType as string}` }); - } + case "update": + { + await secretV2BridgeDAL.updateById( + change.id, + { + skipMultilineEncoding: secretVersion?.skipMultilineEncoding, + version: secretVersion?.version, + type: secretVersion?.type, + key: secretVersion?.key, + reminderNote: secretVersion?.reminderNote, + reminderRepeatDays: secretVersion?.reminderRepeatDays, + encryptedValue: secretVersion?.encryptedValue, + encryptedComment: secretVersion?.encryptedComment, + userId: secretVersion?.userId, + metadata: secretVersion?.metadata + }, + tx + ); + + const newVersion = await secretVersionV2BridgeDAL.create( + { + version: secretVersion.version + 1, + encryptedValue: secretVersion.encryptedValue, + key: secretVersion.key, + encryptedComment: secretVersion.encryptedComment, + skipMultilineEncoding: secretVersion.skipMultilineEncoding, + reminderNote: secretVersion.reminderNote, + reminderRepeatDays: secretVersion.reminderRepeatDays, + userId: secretVersion.userId, + metadata: secretVersion.metadata, + actorType: actorInfo.actorType, + envId: secretVersion.envId, + folderId, + secretId: secretVersion.secretId, + ...(actorInfo.actorType === ActorType.IDENTITY && { identityActorId: actorInfo.actorId }), + ...(actorInfo.actorType === ActorType.USER && { userActorId: actorInfo.actorId }) + }, + tx + ); + + commitChanges.push({ + type: ChangeType.ADD, + secretVersionId: newVersion.id + }); + } + break; + + case "delete": + await secretV2BridgeDAL.deleteById(change.id, tx); + + commitChanges.push({ + type: ChangeType.DELETE, + secretVersionId: change.versionId + }); + break; + + default: + throw new BadRequestError({ message: `Unknown change type: ${change.changeType}` }); } + } - // process folder changes - for (const change of folderChanges) { + return commitChanges; + }; + + /** + * Core function to apply folder state differences + */ + const applyFolderStateDifferencesFn = async ({ + differences, + actorInfo, + folderId, + projectId, + reconstructNewFolders, + reconstructUpToCommit, + step = 0, + tx + }: { + differences: ResourceChange[]; + actorInfo: ActorInfo; + folderId: string; + projectId: string; + reconstructNewFolders: boolean; + reconstructUpToCommit?: string; + step: number; + tx?: Knex; + }): Promise => { + /** + * Process folder changes when applying folder state differences + */ + const processFolderChanges = async ( + changes: ResourceChange[], + folderVersions: Record + ) => { + const commitChanges = []; + + for (const change of changes) { const folderVersion = folderVersions[change.id]; + switch (change.changeType) { case "create": - // Add new folder if (folderVersion) { const newFolder = { id: change.id, @@ -552,45 +689,69 @@ export const folderCommitServiceFactory = ({ tx ); + if (reconstructNewFolders && reconstructUpToCommit && step < 20) { + const subFolderLatestCommit = await folderCommitDAL.findLatestCommitBetween({ + folderId: change.id, + endCommitId: reconstructUpToCommit, + tx + }); + if (subFolderLatestCommit) { + const subFolderDiff = await compareFolderStates({ + targetCommitId: subFolderLatestCommit.id, + tx + }); + if (subFolderDiff?.length > 0) { + await applyFolderStateDifferencesFn({ + differences: subFolderDiff, + actorInfo, + folderId: change.id, + projectId, + reconstructNewFolders, + reconstructUpToCommit, + step: step + 1, + tx + }); + } + } + } + commitChanges.push({ - type: "add", + type: ChangeType.ADD, folderVersionId: newFolderVersion.id }); } break; case "update": - // Update folder to specific version if (change.versionId) { - await folderVersionDAL.findById(change.versionId, tx).then(async (versionDetails) => { - if (versionDetails) { - await folderDAL.updateById( - change.id, - { - parentId: folderId, - envId: versionDetails.envId, - version: (versionDetails.version || 1) + 1, - name: versionDetails.name - }, - tx - ); + const versionDetails = await folderVersionDAL.findById(change.versionId, tx); + if (versionDetails) { + await folderDAL.updateById( + change.id, + { + parentId: folderId, + envId: versionDetails.envId, + version: (versionDetails.version || 1) + 1, + name: versionDetails.name + }, + tx + ); - const newFolderVersion = await folderVersionDAL.create( - { - folderId: change.id, - version: (versionDetails.version || 1) + 1, - name: versionDetails.name, - envId: versionDetails.envId - }, - tx - ); + const newFolderVersion = await folderVersionDAL.create( + { + folderId: change.id, + version: (versionDetails.version || 1) + 1, + name: versionDetails.name, + envId: versionDetails.envId + }, + tx + ); - commitChanges.push({ - type: "add", - folderVersionId: newFolderVersion.id - }); - } - }); + commitChanges.push({ + type: ChangeType.ADD, + folderVersionId: newFolderVersion.id + }); + } } break; @@ -598,78 +759,130 @@ export const folderCommitServiceFactory = ({ await folderDAL.deleteById(change.id, tx); commitChanges.push({ - type: "delete", + type: ChangeType.DELETE, folderVersionId: change.versionId }); break; default: - throw new BadRequestError({ message: `Unknown change type: ${change.changeType as string}` }); + throw new BadRequestError({ message: `Unknown change type: ${change.changeType}` }); } } - await createCommit( - { - actor: { - type: actorInfo.actorType, - metadata: { id: actorInfo.actorId } - }, - message: actorInfo.message || "Rolled back folder state", - folderId, - changes: commitChanges + return commitChanges; + }; + // Group differences by type for more efficient processing + const secretChanges = differences.filter((diff) => diff.type === ResourceType.SECRET); + const folderChanges = differences.filter((diff) => diff.type === ResourceType.FOLDER); + + // Batch fetch necessary data + const secretVersions = await secretVersionV2BridgeDAL.findByIdsWithLatestVersion( + folderId, + secretChanges.map((diff) => diff.id), + secretChanges.map((diff) => diff.versionId) + ); + + const folderVersions = await folderVersionDAL.findByIdsWithLatestVersion( + folderChanges.map((diff) => diff.id), + folderChanges.map((diff) => diff.versionId) + ); + + // Process changes in parallel + const [secretCommitChanges, folderCommitChanges] = await Promise.all([ + processSecretChanges(secretChanges, secretVersions, actorInfo, folderId, tx), + processFolderChanges(folderChanges, folderVersions) + ]); + + // Combine all changes + const allCommitChanges = [...secretCommitChanges, ...folderCommitChanges]; + + // Create a commit with all the changes + await createCommit( + { + actor: { + type: actorInfo.actorType, + metadata: { id: actorInfo.actorId } }, - tx - ); + message: actorInfo.message || "Rolled back folder state", + folderId, + changes: allCommitChanges + }, + tx + ); - result = { - secretChangesCount: secretChanges.length, - folderChangesCount: folderChanges.length, - totalChanges: differences.length - }; - await secretV2BridgeDAL.invalidateSecretCacheByProjectId(projectId); - }); + // Invalidate cache to reflect the changes + await secretV2BridgeDAL.invalidateSecretCacheByProjectId(projectId); - return result; + return { + secretChangesCount: secretChanges.length, + folderChangesCount: folderChanges.length, + totalChanges: differences.length + }; }; - // Retrieve a commit by ID + /** + * Apply folder state differences with transaction handling + */ + const applyFolderStateDifferences = async (params: { + differences: ResourceChange[]; + actorInfo: ActorInfo; + folderId: string; + projectId: string; + reconstructNewFolders: boolean; + reconstructUpToCommit?: string; + tx?: Knex; + }): Promise => { + // If a transaction was provided, use it directly + if (params.tx) { + return applyFolderStateDifferencesFn({ ...params, step: 0 }); + } + + // Otherwise, start a new transaction + return folderCommitDAL.transaction((newTx) => applyFolderStateDifferencesFn({ ...params, tx: newTx, step: 0 })); + }; + + /** + * Retrieve a commit by ID + */ const getCommitById = async (id: string, tx?: Knex) => { return folderCommitDAL.findById(id, tx); }; - // Get all commits for a folder + /** + * Get all commits for a folder + */ const getCommitsByFolderId = async (folderId: string, tx?: Knex) => { return folderCommitDAL.findByFolderId(folderId, tx); }; - // Get changes for a commit + /** + * Get changes for a commit + */ const getCommitChanges = async (commitId: string, tx?: Knex) => { return folderCommitChangesDAL.findByCommitId(commitId, tx); }; - // Get checkpoints for a folder + /** + * Get checkpoints for a folder + */ const getCheckpointsByFolderId = async (folderId: string, limit?: number, tx?: Knex) => { return folderCheckpointDAL.findByFolderId(folderId, limit, tx); }; - // Get the latest checkpoint for a folder + /** + * Get the latest checkpoint for a folder + */ const getLatestCheckpoint = async (folderId: string, tx?: Knex) => { return folderCheckpointDAL.findLatestByFolderId(folderId, tx); }; - // Get tree checkpoints for a project - const getTreeCheckpointsByProjectId = async (projectId: string, limit?: number, tx?: Knex) => { - return folderTreeCheckpointDAL.findByProjectId(projectId, limit, tx); - }; - - // Get the latest tree checkpoint for a project - const getLatestTreeCheckpoint = async (projectId: string, tx?: Knex) => { - return folderTreeCheckpointDAL.findLatestByProjectId(projectId, tx); - }; - + /** + * Initialize a folder with its current state + */ const initializeFolder = async (folderId: string, tx?: Knex) => { const folderResources = await getFolderResources(folderId, tx); - const changes = folderResources.map((resource) => ({ type: "add", ...resource })); + const changes = folderResources.map((resource) => ({ type: ChangeType.ADD, ...resource })); + if (changes.length > 0) { const newCommit = await createCommit( { @@ -686,50 +899,225 @@ export const folderCommitServiceFactory = ({ } }; - function sortFoldersByHierarchy(folders: TSecretFolders[]) { + /** + * Sort folders by hierarchy (parents before children) + */ + const sortFoldersByHierarchy = (folders: TSecretFolders[]) => { // Create a map for quick lookup of children by parent ID - const childrenMap: Map = new Map(); + const childrenMap = new Map(); + + // Set of all folder IDs + const allFolderIds = new Set(); + + // Build the set of all folder IDs folders.forEach((folder) => { - const { parentId } = folder; - if (!childrenMap.has(parentId || null)) { - childrenMap.set(parentId || null, []); + if (folder.id) { + allFolderIds.add(folder.id); } - childrenMap.get(parentId || null)?.push(folder); }); - // Start with root folders (null parentId) - const result = []; - const rootFolders = childrenMap.get(null) || []; + // Group folders by their parentId + folders.forEach((folder) => { + if (folder.parentId) { + const children = childrenMap.get(folder.parentId) || []; + children.push(folder); + childrenMap.set(folder.parentId, children); + } + }); + + // Find root folders - those with no parentId or with a parentId that doesn't exist + const rootFolders = folders.filter((folder) => !folder.parentId || !allFolderIds.has(folder.parentId)); // Process each level of the hierarchy + const result = []; let currentLevel = rootFolders; - result.push(...currentLevel); while (currentLevel.length > 0) { - const nextLevel = []; + result.push(...currentLevel); + const nextLevel = []; for (const folder of currentLevel) { - const children = childrenMap.get(folder.id) || []; - nextLevel.push(...children); + if (folder.id) { + const children = childrenMap.get(folder.id) || []; + nextLevel.push(...children); + } } - result.push(...nextLevel); currentLevel = nextLevel; } return result; - } + }; + /** + * Create a checkpoint for a folder tree + */ + const createFolderTreeCheckpoint = async (envId: string, folderCommitId?: string, tx?: Knex) => { + let latestCommitId = folderCommitId; + const latestTreeCheckpoint = await folderTreeCheckpointDAL.findLatestByEnvId(envId, tx); + + if (!latestCommitId) { + const latestCommit = await folderCommitDAL.findLatestEnvCommit(envId, tx); + if (!latestCommit) { + logger.info(`createFolderTreeCheckpoint - Latest commit ID not found for envId ${envId}`); + return; + } + latestCommitId = latestCommit.id; + } + + if (latestTreeCheckpoint) { + const commitsSinceLastCheckpoint = await folderCommitDAL.getEnvNumberOfCommitsSince( + envId, + latestTreeCheckpoint.folderCommitId, + tx + ); + if (commitsSinceLastCheckpoint < Number(appCfg.PIT_TREE_CHECKPOINT_WINDOW)) { + logger.info( + `createFolderTreeCheckpoint - Commits since last checkpoint ${commitsSinceLastCheckpoint} is less than ${appCfg.PIT_TREE_CHECKPOINT_WINDOW}` + ); + return; + } + } + + const folders = await folderDAL.findByEnvId(envId, tx); + const sortedFolders = sortFoldersByHierarchy(folders); + const filteredFoldersIds = sortedFolders.filter((folder) => !folder.isReserved).map((folder) => folder.id); + const folderCommits = await folderCommitDAL.findMultipleLatestCommits(filteredFoldersIds, tx); + const folderTreeCheckpoint = await folderTreeCheckpointDAL.create( + { + folderCommitId: latestCommitId + }, + tx + ); + await folderTreeCheckpointResourcesDAL.insertMany( + folderCommits.map((folderCommit) => ({ + folderTreeCheckpointId: folderTreeCheckpoint.id, + folderId: folderCommit.folderId, + folderCommitId: folderCommit.id + })), + tx + ); + }; + + /** + * Initialize a project with its current state + */ const initializeProject = async (projectId: string, tx?: Knex) => { const project = await projectDAL.findById(projectId, tx); if (!project) { throw new NotFoundError({ message: `Project with ID ${projectId} not found` }); } + const folders = await folderDAL.findByProjectId(projectId, tx); const sortedFolders = sortFoldersByHierarchy(folders); + await Promise.all(sortedFolders.map((folder) => initializeFolder(folder.id, tx))); + + const envIds = [...new Set(folders.map((folder) => folder.envId))]; + await Promise.all( + envIds.map(async (envId) => { + await createFolderTreeCheckpoint(envId, undefined, tx); + }) + ); }; + /** + * Roll back a folder tree to a specific commit + */ + const deepRollbackFolder = async ( + targetCommitId: string, + envId: string, + actorId: string, + actorType: ActorType, + projectId: string, + tx?: Knex + ) => { + const targetCommit = await folderCommitDAL.findById(targetCommitId, tx); + if (!targetCommit) { + throw new NotFoundError({ message: `No commit found for commit ID ${targetCommitId}` }); + } + + const checkpoint = await folderTreeCheckpointDAL.findNearestCheckpoint(targetCommitId, envId, tx); + if (!checkpoint) { + throw new NotFoundError({ message: `No checkpoint found for commit ID ${targetCommitId}` }); + } + + const folderCheckpointCommits = await folderTreeCheckpointResourcesDAL.findByTreeCheckpointId(checkpoint.id, tx); + const folderCommits = await folderCommitDAL.findAllCommitsBetween({ + envId, + endCommitId: targetCommit.commitId.toString(), + startCommitId: checkpoint.commitId.toString(), + tx + }); + + // Group commits by folderId and keep only the latest + const folderGroups = new Map(); + + if (folderCheckpointCommits && folderCheckpointCommits.length > 0) { + for (const commit of folderCheckpointCommits) { + folderGroups.set(commit.folderId, { + createdAt: commit.createdAt, + id: commit.folderCommitId + }); + } + } + + if (folderCommits && folderCommits.length > 0) { + for (const commit of folderCommits) { + const { folderId, createdAt, id } = commit; + const existingCommit = folderGroups.get(folderId); + + if (!existingCommit || createdAt.getTime() > existingCommit.createdAt.getTime()) { + folderGroups.set(folderId, { createdAt, id }); + } + } + } + + const folderDiffs = new Map(); + + // Process each folder to determine differences + await Promise.all( + Array.from(folderGroups.entries()).map(async ([folderId, commit]) => { + const latestFolderCommit = await folderCommitDAL.findLatestCommit(folderId, tx); + if (latestFolderCommit && latestFolderCommit.id !== commit.id) { + const diff = await compareFolderStates({ + currentCommitId: latestFolderCommit.id, + targetCommitId: commit.id, + tx + }); + if (diff?.length > 0) { + folderDiffs.set(folderId, diff); + } + } + }) + ); + + // Apply changes in hierarchical order + const folderIds = Array.from(folderDiffs.keys()); + const folders = await folderDAL.findFoldersByRootAndIds({ rootId: targetCommit.folderId, folderIds }, tx); + const sortedFolders = sortFoldersByHierarchy(folders); + + for (const folder of sortedFolders) { + const diff = folderDiffs.get(folder.id); + if (diff) { + await applyFolderStateDifferences({ + differences: diff, + actorInfo: { + actorType, + actorId, + message: "Deep rollback" + }, + folderId: folder.id, + projectId, + reconstructNewFolders: true, + reconstructUpToCommit: targetCommit.commitId.toString(), + tx + }); + } + } + }; + + // Return the public interface return { createCommit, addCommitChange, @@ -738,13 +1126,13 @@ export const folderCommitServiceFactory = ({ getCommitChanges, getCheckpointsByFolderId, getLatestCheckpoint, - getTreeCheckpointsByProjectId, - getLatestTreeCheckpoint, initializeFolder, initializeProject, createFolderCheckpoint, compareFolderStates, - applyFolderStateDifferences + applyFolderStateDifferences, + createFolderTreeCheckpoint, + deepRollbackFolder }; }; diff --git a/backend/src/services/folder-tree-checkpoint-resources/folder-tree-checkpoint-resources-dal.ts b/backend/src/services/folder-tree-checkpoint-resources/folder-tree-checkpoint-resources-dal.ts index 8813784fd..5d3fac9c2 100644 --- a/backend/src/services/folder-tree-checkpoint-resources/folder-tree-checkpoint-resources-dal.ts +++ b/backend/src/services/folder-tree-checkpoint-resources/folder-tree-checkpoint-resources-dal.ts @@ -1,27 +1,14 @@ import { Knex } from "knex"; import { TDbClient } from "@app/db"; -import { - TableName, - TFolderTreeCheckpointResources, - TFolderTreeCheckpoints, - TProjectEnvironments, - TSecretFolders -} from "@app/db/schemas"; +import { TableName, TFolderTreeCheckpointResources } from "@app/db/schemas"; import { DatabaseError } from "@app/lib/errors"; -import { ormify, selectAllTableCols } from "@app/lib/knex"; +import { buildFindFilter, ormify, selectAllTableCols } from "@app/lib/knex"; export type TFolderTreeCheckpointResourcesDALFactory = ReturnType; -type ResourceWithCheckpointInfo = TFolderTreeCheckpointResources & { - folderCommitId: string; -}; - -type ResourceWithFolderInfo = TFolderTreeCheckpointResources & { - name: string; - parentId?: string | null; - slug: string; - envName: string; +type TFolderTreeCheckpointResourcesWithCommitId = TFolderTreeCheckpointResources & { + commitId: number; }; export const folderTreeCheckpointResourcesDALFactory = (db: TDbClient) => { @@ -30,95 +17,28 @@ export const folderTreeCheckpointResourcesDALFactory = (db: TDbClient) => { const findByTreeCheckpointId = async ( folderTreeCheckpointId: string, tx?: Knex - ): Promise => { + ): Promise => { try { const docs = await (tx || db.replicaNode())( TableName.FolderTreeCheckpointResources ) - .where({ folderTreeCheckpointId }) - .select(selectAllTableCols(TableName.FolderTreeCheckpointResources)); + .join( + TableName.FolderCommit, + `${TableName.FolderTreeCheckpointResources}.folderCommitId`, + `${TableName.FolderCommit}.id` + ) + // eslint-disable-next-line @typescript-eslint/no-misused-promises + .where(buildFindFilter({ folderTreeCheckpointId }, TableName.FolderTreeCheckpointResources)) + .select(selectAllTableCols(TableName.FolderTreeCheckpointResources)) + .select(db.ref("commitId").withSchema(TableName.FolderCommit).as("commitId")); return docs; } catch (error) { throw new DatabaseError({ error, name: "FindByTreeCheckpointId" }); } }; - const findByFolderId = async (folderId: string, tx?: Knex): Promise => { - try { - const docs = await (tx || db.replicaNode())< - TFolderTreeCheckpointResources & Pick - >(TableName.FolderTreeCheckpointResources) - .where({ folderId }) - .select(selectAllTableCols(TableName.FolderTreeCheckpointResources)) - .join( - TableName.FolderTreeCheckpoint, - `${TableName.FolderTreeCheckpointResources}.folderTreeCheckpointId`, - `${TableName.FolderTreeCheckpoint}.id` - ) - .select( - db.ref("folderCommitId").withSchema(TableName.FolderTreeCheckpoint), - db.ref("createdAt").withSchema(TableName.FolderTreeCheckpoint) - ); - return docs; - } catch (error) { - throw new DatabaseError({ error, name: "FindByFolderId" }); - } - }; - - const findByFolderCommitId = async (folderCommitId: string, tx?: Knex): Promise => { - try { - const docs = await (tx || db.replicaNode())< - TFolderTreeCheckpointResources & Pick - >(TableName.FolderTreeCheckpointResources) - .where({ folderCommitId }) - .select(selectAllTableCols(TableName.FolderTreeCheckpointResources)) - .join( - TableName.FolderTreeCheckpoint, - `${TableName.FolderTreeCheckpointResources}.folderTreeCheckpointId`, - `${TableName.FolderTreeCheckpoint}.id` - ) - .select(db.ref("createdAt").withSchema(TableName.FolderTreeCheckpoint)); - return docs; - } catch (error) { - throw new DatabaseError({ error, name: "FindByFolderCommitId" }); - } - }; - - const findFoldersInTreeCheckpoint = async ( - folderTreeCheckpointId: string, - tx?: Knex - ): Promise => { - try { - const docs = await (tx || db.replicaNode())< - TFolderTreeCheckpointResources & - Pick & - Pick & { envName: string } - >(TableName.FolderTreeCheckpointResources) - .where({ folderTreeCheckpointId }) - .select(selectAllTableCols(TableName.FolderTreeCheckpointResources)) - .join( - TableName.SecretFolder, - `${TableName.FolderTreeCheckpointResources}.folderId`, - `${TableName.SecretFolder}.id` - ) - .join(TableName.Environment, `${TableName.SecretFolder}.envId`, `${TableName.Environment}.id`) - .select( - db.ref("name").withSchema(TableName.SecretFolder), - db.ref("parentId").withSchema(TableName.SecretFolder), - db.ref("slug").withSchema(TableName.Environment), - db.ref("name").withSchema(TableName.Environment).as("envName") - ); - return docs; - } catch (error) { - throw new DatabaseError({ error, name: "FindFoldersInTreeCheckpoint" }); - } - }; - return { ...folderTreeCheckpointResourcesOrm, - findByTreeCheckpointId, - findByFolderId, - findByFolderCommitId, - findFoldersInTreeCheckpoint + findByTreeCheckpointId }; }; diff --git a/backend/src/services/folder-tree-checkpoint/folder-tree-checkpoint-dal.ts b/backend/src/services/folder-tree-checkpoint/folder-tree-checkpoint-dal.ts index 4e9381f9c..7589c4b49 100644 --- a/backend/src/services/folder-tree-checkpoint/folder-tree-checkpoint-dal.ts +++ b/backend/src/services/folder-tree-checkpoint/folder-tree-checkpoint-dal.ts @@ -8,11 +8,7 @@ import { ormify, selectAllTableCols } from "@app/lib/knex"; export type TFolderTreeCheckpointDALFactory = ReturnType; type TreeCheckpointWithCommitInfo = TFolderTreeCheckpoints & { - actorMetadata: unknown; - actorType: string; - message?: string | null; - commitDate: Date; - folderId: string; + commitId: number; }; export const folderTreeCheckpointDALFactory = (db: TDbClient) => { @@ -30,82 +26,61 @@ export const folderTreeCheckpointDALFactory = (db: TDbClient) => { } }; - const findByProjectId = async ( - projectId: string, - limit?: number, - tx?: Knex - ): Promise => { - try { - const query = (tx || db.replicaNode())< - TFolderTreeCheckpoints & - Pick & { commitDate: Date } - >(TableName.FolderTreeCheckpoint) - .join( - TableName.FolderCommit, - `${TableName.FolderTreeCheckpoint}.folderCommitId`, - `${TableName.FolderCommit}.id` - ) - .join(TableName.SecretFolder, `${TableName.FolderCommit}.folderId`, `${TableName.SecretFolder}.id`) - .join(TableName.Environment, `${TableName.SecretFolder}.envId`, `${TableName.Environment}.id`) - .where({ projectId }) - .select(selectAllTableCols(TableName.FolderTreeCheckpoint)) - .select( - db.ref("actorMetadata").withSchema(TableName.FolderCommit), - db.ref("actorType").withSchema(TableName.FolderCommit), - db.ref("message").withSchema(TableName.FolderCommit), - db.ref("createdAt").withSchema(TableName.FolderCommit).as("commitDate"), - db.ref("folderId").withSchema(TableName.FolderCommit) - ) - .orderBy(`${TableName.FolderTreeCheckpoint}.createdAt`, "desc"); - - if (limit) { - void query.limit(limit); - } - - const docs = await query; - return docs; - } catch (error) { - throw new DatabaseError({ error, name: "FindByProjectId" }); - } - }; - - const findLatestByProjectId = async ( - projectId: string, + const findNearestCheckpoint = async ( + folderCommitId: string, + envId: string, tx?: Knex ): Promise => { try { - const doc = await (tx || db.replicaNode())< - TFolderTreeCheckpoints & - Pick & { commitDate: Date } - >(TableName.FolderTreeCheckpoint) - .join( + const targetCommit = await (tx || db.replicaNode())(TableName.FolderCommit) + .where({ id: folderCommitId }) + .select("id", "commitId", "folderId") + .first(); + + if (!targetCommit) { + return undefined; + } + + const nearestCheckpoint = await (tx || db.replicaNode())(TableName.FolderTreeCheckpoint) + .join( TableName.FolderCommit, `${TableName.FolderTreeCheckpoint}.folderCommitId`, `${TableName.FolderCommit}.id` ) - .join(TableName.SecretFolder, `${TableName.FolderCommit}.folderId`, `${TableName.SecretFolder}.id`) - .join(TableName.Environment, `${TableName.SecretFolder}.envId`, `${TableName.Environment}.id`) - .where({ projectId }) + .where(`${TableName.FolderCommit}.commitId`, "<=", targetCommit.commitId.toString()) + .where(`${TableName.FolderCommit}.envId`, envId) .select(selectAllTableCols(TableName.FolderTreeCheckpoint)) - .select( - db.ref("actorMetadata").withSchema(TableName.FolderCommit), - db.ref("actorType").withSchema(TableName.FolderCommit), - db.ref("message").withSchema(TableName.FolderCommit), - db.ref("createdAt").withSchema(TableName.FolderCommit).as("commitDate"), - db.ref("folderId").withSchema(TableName.FolderCommit) + .select(db.ref("commitId").withSchema(TableName.FolderCommit)) + .orderBy(`${TableName.FolderCommit}.commitId`, "desc") + .first(); + + return nearestCheckpoint; + } catch (error) { + throw new DatabaseError({ error, name: "FindNearestCheckpoint" }); + } + }; + + const findLatestByEnvId = async (envId: string, tx?: Knex): Promise => { + try { + const doc = await (tx || db.replicaNode())(TableName.FolderTreeCheckpoint) + .join( + TableName.FolderCommit, + `${TableName.FolderTreeCheckpoint}.folderCommitId`, + `${TableName.FolderCommit}.id` ) + .where(`${TableName.FolderCommit}.envId`, envId) .orderBy(`${TableName.FolderTreeCheckpoint}.createdAt`, "desc") .first(); return doc; } catch (error) { - throw new DatabaseError({ error, name: "FindLatestByProjectId" }); + throw new DatabaseError({ error, name: "FindLatestByFolderId" }); } }; return { ...folderTreeCheckpointOrm, findByCommitId, - findByProjectId, - findLatestByProjectId + findNearestCheckpoint, + findLatestByEnvId }; }; diff --git a/backend/src/services/secret-folder/secret-folder-dal.ts b/backend/src/services/secret-folder/secret-folder-dal.ts index d6d5c4aad..5e5809ba5 100644 --- a/backend/src/services/secret-folder/secret-folder-dal.ts +++ b/backend/src/services/secret-folder/secret-folder-dal.ts @@ -488,6 +488,51 @@ export const secretFolderDALFactory = (db: TDbClient) => { } }; + const findFoldersByRootAndIds = async ({ rootId, folderIds }: { rootId: string; folderIds: string[] }, tx?: Knex) => { + try { + // First, get all descendant folders of rootId + const descendants = await (tx || db.replicaNode()) + .withRecursive("descendants", (qb) => + qb + .select( + selectAllTableCols(TableName.SecretFolder), + db.raw("0 as depth"), + db.raw(`'/' as path`), + db.ref(`${TableName.Environment}.slug`).as("environment") + ) + .from(TableName.SecretFolder) + .join(TableName.Environment, `${TableName.SecretFolder}.envId`, `${TableName.Environment}.id`) + .where(`${TableName.SecretFolder}.id`, rootId) + .union((un) => { + void un + .select( + selectAllTableCols(TableName.SecretFolder), + db.raw("descendants.depth + 1 as depth"), + db.raw( + `CONCAT( + CASE WHEN descendants.path = '/' THEN '' ELSE descendants.path END, + CASE WHEN ${TableName.SecretFolder}."parentId" is NULL THEN '' ELSE CONCAT('/', secret_folders.name) END + )` + ), + db.ref("descendants.environment") + ) + .from(TableName.SecretFolder) + .where(`${TableName.SecretFolder}.isReserved`, false) + .join("descendants", `${TableName.SecretFolder}.parentId`, "descendants.id"); + }) + ) + .select<(TSecretFolders & { path: string; depth: number; environment: string })[]>("*") + .from("descendants") + .whereIn(`id`, folderIds) + .orderBy("depth") + .orderBy(`name`); + + return descendants; + } catch (error) { + throw new DatabaseError({ error, name: "FindFoldersByRootAndIds" }); + } + }; + const findByParentId = async (parentId: string, tx?: Knex) => { try { const folders = await (tx || db.replicaNode())(TableName.SecretFolder) @@ -499,6 +544,17 @@ export const secretFolderDALFactory = (db: TDbClient) => { } }; + const findByEnvId = async (envId: string, tx?: Knex) => { + try { + const folders = await (tx || db.replicaNode())(TableName.SecretFolder) + .where({ envId }) + .select(selectAllTableCols(TableName.SecretFolder)); + return folders; + } catch (error) { + throw new DatabaseError({ error, name: "findByEnvId" }); + } + }; + return { ...secretFolderOrm, update, @@ -511,6 +567,8 @@ export const secretFolderDALFactory = (db: TDbClient) => { findByProjectId, findByMultiEnv, findByEnvsDeep, - findByParentId + findByParentId, + findByEnvId, + findFoldersByRootAndIds }; };