PIT: rework of init migration

This commit is contained in:
carlosmonastyrski
2025-05-29 16:44:20 -03:00
parent d1d5dd29c6
commit a70aff5f31
2 changed files with 237 additions and 48 deletions
@@ -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<string, TSecretFolders[]>();
// Set of all folder IDs
const allFolderIds = new Set<string>();
// 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<void> {
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<string, string> = {};
// 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<string, string> = {};
const newCommitsMapInverted: Record<string, string> = {};
const newCheckpointsMap: Record<string, string> = {};
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<string, string> = {};
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<void> {
@@ -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,