diff --git a/backend/src/db/migrations/20250507185010_pit-projects-commits-initialization.ts b/backend/src/db/migrations/20250507185010_pit-projects-commits-initialization.ts index 459727035..3fa57bb98 100644 --- a/backend/src/db/migrations/20250507185010_pit-projects-commits-initialization.ts +++ b/backend/src/db/migrations/20250507185010_pit-projects-commits-initialization.ts @@ -1,23 +1,231 @@ import { Knex } from "knex"; import { inMemoryKeyStore } from "@app/keystore/memory"; +import { chunkArray } from "@app/lib/fn"; +import { selectAllTableCols } from "@app/lib/knex"; +import { logger } from "@app/lib/logger"; -import { ProjectType, TableName } from "../schemas"; +import { + ProjectType, + TableName, + TFolderCheckpoints, + TFolderCommits, + TFolderTreeCheckpoints, + TSecretFolders +} from "../schemas"; import { getMigrationEnvConfig } from "./utils/env-config"; import { getMigrationPITServices } from "./utils/services"; +const sortFoldersByHierarchy = (folders: TSecretFolders[]) => { + // Create a map for quick lookup of children by parent ID + const childrenMap = new Map(); + + // Set of all folder IDs + const allFolderIds = new Set(); + + // Build the set of all folder IDs + folders.forEach((folder) => { + if (folder.id) { + allFolderIds.add(folder.id); + } + }); + + // 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; + + 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.reverse(); +}; + export async function up(knex: Knex): Promise { + logger.info("Initializing folder commits"); const hasFolderCommitTable = await knex.schema.hasTable(TableName.FolderCommit); if (hasFolderCommitTable) { const keyStore = inMemoryKeyStore(); const envConfig = getMigrationEnvConfig(); const { folderCommitService } = await getMigrationPITServices({ db: knex, keyStore, envConfig }); - const projects = await knex(TableName.Project).where({ version: 3, type: ProjectType.SecretManager }).select("id"); - for (const project of projects) { + + // Get Projects to Initialize + const projects = await knex(TableName.Project) + .where(`${TableName.Project}.version`, 3) + .where(`${TableName.Project}.type`, ProjectType.SecretManager) + .select(selectAllTableCols(TableName.Project)); + logger.info(`Found ${projects.length} projects to initialize`); + + // Process Projects in batches of 100 + const batches = chunkArray(projects, 100); + let i = 0; + for (const batch of batches) { + i += 1; + logger.info(`Processing project batch ${i} of ${batches.length}`); + const foldersCommitsList = []; + + const rootFoldersMap: Record = {}; + + // Get All Folders for the Project // eslint-disable-next-line no-await-in-loop - await folderCommitService.initializeProject(project.id, knex); + const folders = await knex(TableName.SecretFolder) + .join(TableName.Environment, `${TableName.SecretFolder}.envId`, `${TableName.Environment}.id`) + .whereIn( + `${TableName.Environment}.projectId`, + batch.map((project) => project.id) + ) + .select(selectAllTableCols(TableName.SecretFolder)); + logger.info(`Found ${folders.length} folders to initialize in project batch ${i} of ${batches.length}`); + + // Sort Folders by Hierarchy (parents before nested folders) + const sortedFolders = sortFoldersByHierarchy(folders); + + // Get folder commit changes + for (const folder of sortedFolders) { + // eslint-disable-next-line no-await-in-loop + const folderCommit = await folderCommitService.getFolderInitialChanges(folder.id, folder.envId, knex); + if (folderCommit.commit && folderCommit.changes) { + foldersCommitsList.push(folderCommit); + if (!folder.parentId) { + rootFoldersMap[folder.id] = folder.envId; + } + } + } + logger.info(`Retrieved folder changes for project batch ${i} of ${batches.length}`); + + // Insert New Commits in batches of 9000 + const newCommits = foldersCommitsList.map((folderCommit) => folderCommit.commit); + const commitBatches = chunkArray(newCommits, 9000); + + let j = 0; + for (const commitBatch of commitBatches) { + j += 1; + logger.info(`Inserting folder commits - batch ${j} of ${commitBatches.length}`); + // Create folder commit + // eslint-disable-next-line no-await-in-loop + const newCommitsInserted = (await knex + .batchInsert(TableName.FolderCommit, commitBatch) + .returning("*")) as TFolderCommits[]; + + logger.info(`Finished inserting folder commits - batch ${j} of ${commitBatches.length}`); + + const newCommitsMap: Record = {}; + const newCommitsMapInverted: Record = {}; + const newCheckpointsMap: Record = {}; + newCommitsInserted.forEach((commit) => { + newCommitsMap[commit.folderId] = commit.id; + newCommitsMapInverted[commit.id] = commit.folderId; + }); + + // Create folder checkpoints + // eslint-disable-next-line no-await-in-loop + const newCheckpoints = (await knex + .batchInsert( + TableName.FolderCheckpoint, + Object.values(newCommitsMap).map((commitId) => ({ + folderCommitId: commitId + })) + ) + .returning("*")) as TFolderCheckpoints[]; + + logger.info(`Finished inserting folder checkpoints - batch ${j} of ${commitBatches.length}`); + + newCheckpoints.forEach((checkpoint) => { + newCheckpointsMap[newCommitsMapInverted[checkpoint.folderCommitId]] = checkpoint.id; + }); + + // Create folder commit changes + // eslint-disable-next-line no-await-in-loop + await knex.batchInsert( + TableName.FolderCommitChanges, + foldersCommitsList + .map((folderCommit) => folderCommit.changes) + .flat() + .map((change) => ({ + folderCommitId: newCommitsMap[change.folderId], + changeType: change.changeType, + secretVersionId: change.secretVersionId, + folderVersionId: change.folderVersionId, + isUpdate: false + })) + ); + + logger.info(`Finished inserting folder commit changes - batch ${j} of ${commitBatches.length}`); + + // Create folder checkpoint resources + // eslint-disable-next-line no-await-in-loop + await knex.batchInsert( + TableName.FolderCheckpointResources, + foldersCommitsList + .map((folderCommit) => folderCommit.changes) + .flat() + .map((change) => ({ + folderCheckpointId: newCheckpointsMap[change.folderId], + folderVersionId: change.folderVersionId, + secretVersionId: change.secretVersionId + })) + ); + + logger.info(`Finished inserting folder checkpoint resources - batch ${j} of ${commitBatches.length}`); + + // Create Folder Tree Checkpoint + // eslint-disable-next-line no-await-in-loop + const newTreeCheckpoints = (await knex + .batchInsert( + TableName.FolderTreeCheckpoint, + Object.keys(rootFoldersMap).map((folderId) => ({ + folderCommitId: newCommitsMap[folderId] + })) + ) + .returning("*")) as TFolderTreeCheckpoints[]; + + logger.info(`Finished inserting folder tree checkpoints - batch ${j} of ${commitBatches.length}`); + + const newTreeCheckpointsMap: Record = {}; + newTreeCheckpoints.forEach((checkpoint) => { + newTreeCheckpointsMap[rootFoldersMap[newCommitsMapInverted[checkpoint.folderCommitId]]] = checkpoint.id; + }); + + // Create Folder Tree Checkpoint Resources + // eslint-disable-next-line no-await-in-loop + await knex + .batchInsert( + TableName.FolderTreeCheckpointResources, + newCommitsInserted.map((folderCommit) => ({ + folderTreeCheckpointId: newTreeCheckpointsMap[folderCommit.envId], + folderId: folderCommit.folderId, + folderCommitId: folderCommit.id + })) + ) + .returning("*"); + + logger.info(`Finished inserting folder tree checkpoint resources - batch ${j} of ${commitBatches.length}`); + } } } + logger.info("Folder commits initialized"); } export async function down(knex: Knex): Promise { diff --git a/backend/src/services/folder-commit/folder-commit-service.ts b/backend/src/services/folder-commit/folder-commit-service.ts index 8fe89e8c2..aac3804aa 100644 --- a/backend/src/services/folder-commit/folder-commit-service.ts +++ b/backend/src/services/folder-commit/folder-commit-service.ts @@ -240,12 +240,19 @@ export const folderCommitServiceFactory = ({ 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 }))); + resources.push( + ...Object.values(folderVersions).map((folderVersion) => ({ + folderVersionId: folderVersion.id, + secretVersionId: undefined + })) + ); } const secretVersions = await secretVersionV2BridgeDAL.findLatestVersionByFolderId(folderId, tx); if (secretVersions.length > 0) { - resources.push(...secretVersions.map((secretVersion) => ({ secretVersionId: secretVersion.id }))); + resources.push( + ...secretVersions.map((secretVersion) => ({ secretVersionId: secretVersion.id, folderVersionId: undefined })) + ); } return resources; @@ -1574,27 +1581,29 @@ export const folderCommitServiceFactory = ({ /** * Initialize a folder with its current state */ - const initializeFolder = async (folderId: string, tx?: Knex) => { + const getFolderInitialChanges = async (folderId: string, envId: string, tx?: Knex) => { const folderResources = await getFolderResources(folderId, tx); const changes = folderResources.map((resource) => ({ type: ChangeType.ADD, ...resource })); if (changes.length > 0) { - const newCommit = await createCommit( - { - actor: { - type: ActorType.PLATFORM - }, + return { + commit: { + actorMetadata: {}, + actorType: ActorType.PLATFORM, message: "Initialized folder", folderId, - changes, - omitIgnoreFilter: true + envId }, - tx - ); - if (newCommit) { - await createFolderCheckpoint({ folderId, folderCommitId: newCommit.id, force: true, tx }); - } + changes: changes.map((change) => ({ + folderId, + changeType: change.type, + secretVersionId: change.secretVersionId, + folderVersionId: change.folderVersionId, + isUpdate: false + })) + }; } + return {}; }; /** @@ -1697,33 +1706,6 @@ export const folderCommitServiceFactory = ({ ); }; - /** - * 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); - - const batchSize = 5; - const chunks = chunkArray(sortedFolders, batchSize); - - for (const chunk of chunks) { - await Promise.all(chunk.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); - }) - ); - }; - const addNestedFolderChanges = async ({ changes, beforeCommit, @@ -2173,8 +2155,7 @@ export const folderCommitServiceFactory = ({ getCommitChanges, getCheckpointsByFolderId, getLatestCheckpoint, - initializeFolder, - initializeProject, + getFolderInitialChanges, createFolderCheckpoint, compareFolderStates, applyFolderStateDifferences,