diff --git a/pg-migrator/src/index.ts b/pg-migrator/src/index.ts index ba43920be..b9ecf85ba 100644 --- a/pg-migrator/src/index.ts +++ b/pg-migrator/src/index.ts @@ -5,6 +5,8 @@ import knex, { Knex } from "knex"; import path from "path"; import { Level } from "level"; import { packRules } from "@casl/ability/extra"; +import slugify from '@sindresorhus/slugify'; + import { APIKeyData, BackupPrivateKey, @@ -57,6 +59,8 @@ enum SecretEncryptionAlgo { AES_256_GCM = "aes-256-gcm", } +const ENV_SLUG_LENGTH = 15; + enum SecretKeyEncoding { UTF8 = "utf8", BASE64 = "base64", @@ -69,6 +73,12 @@ const getFolderVersionKey = (folderId: string, version: number) => const projectKv = kdb.sublevel(TableName.Project); +const migrationCheckPointsKv = kdb.sublevel("CHECK-POINTS"); + +export const truncateAndSlugify = (slug: string) : string => { + return slugify(slug.slice(0, ENV_SLUG_LENGTH)) +} + /** * Sorts an array of items into groups. The return value is a map where the keys are * the group ids the given getGroupId function produced and the value is an array of @@ -118,10 +128,24 @@ const migrateCollection = async < mongooseCollection: Model; filter?: Record; }) => { - const mongooseDoc: T[] = []; - const pgDoc: Tables[K]["base"][] = []; - console.log("Starting migration of ", mongooseCollection.modelName); + // check ones that have already been migrated + const migrationCheckPointsKvRes = await migrationCheckPointsKv.get(postgresTableName).catch(() => null) + if (migrationCheckPointsKvRes) { + console.log(`Skipping Postgres table '${postgresTableName}' because of check point`) + return + } else { + await db(postgresTableName).del() + console.log(`Dropped table named '${postgresTableName}' to start fresh from this point on`) + } + + // clear table to start fresh from this check point + + + const mongooseDoc: T[] = []; + const pgDoc: Tables[K]["base"][] = []; // pre processed data ready to be inserted into PSQL + + console.log("Starting migration of ", mongooseCollection.modelName, " Postgres table name:", postgresTableName); const totalMongoCount = await mongooseCollection.countDocuments(); console.log("Total documents", totalMongoCount); console.log("Total batches", Math.ceil(totalMongoCount / 1000)); @@ -129,9 +153,7 @@ const migrateCollection = async < for await (const doc of mongooseCollection.find(filter || {}).cursor({ batchSize: 100 })) { mongooseDoc.push(doc); - const preProcessedData = await preProcessing( - doc.toObject({ virtuals: true }), - ); + const preProcessedData = await preProcessing(doc.toObject({ virtuals: true })); if (preProcessedData) { if (Array.isArray(preProcessedData)) { pgDoc.push( @@ -181,12 +203,17 @@ const migrateCollection = async < pgDoc.splice(0, pgDoc.length); } - console.log("Finished migration of ", mongooseCollection.modelName); + migrationCheckPointsKv.put(postgresTableName, "done") + + console.log("Finished migration of ", mongooseCollection.modelName, " Postgres table name:", postgresTableName); }; const main = async () => { try { dotenv.config(); + + process.env.START_FRESH = "true"; + const prompt = promptSync({ sigint: true }); let mongodb_url = process.env.MONGO_DB_URL; @@ -215,14 +242,18 @@ const main = async () => { console.log("Connected successfully to postgres"); await db.raw("select 1+1 as result"); - console.log("Starting rolling back to latest, comment this out later"); - await db.migrate.rollback({}, true); - await kdb.clear(); - console.log("Rolling back completed"); + if (process.env.START_FRESH === "true") { + await migrationCheckPointsKv.clear() - console.log("Executing migration"); - await db.migrate.latest(); - console.log("Completed migration"); + console.log("Starting rolling back to latest, comment this out later"); + // await db.migrate.rollback({}, true); + await kdb.clear(); + console.log("Rolling back completed"); + + console.log("Executing migration"); + await db.migrate.latest(); + console.log("Completed migration"); + } const userKv = kdb.sublevel(TableName.Users); await migrateCollection({ @@ -232,10 +263,8 @@ const main = async () => { returnKeys: ["id", "email"], preProcessing: async (doc) => { const id = uuidV4(); - const userKeyExists = await userKv.get(id).catch(()=> null) - if (!userKeyExists) return - userKv.put(doc.id, id); + await userKv.put(doc.id.toString(), id); return { id, @@ -351,7 +380,7 @@ const main = async () => { name: doc.name, orgId, description: doc.description, - slug: doc.slug, + slug: truncateAndSlugify(doc.slug), permissions: doc.permissions ? JSON.stringify(packRules(doc.permissions as any)) : null, @@ -517,6 +546,12 @@ const main = async () => { return envKv.sublevel(environment); }; + const checkIfFolderIsDangling = async (projectId: string, env_slug: string, folderId: string) => { + const kv = getFolderKv(projectId, truncateAndSlugify(env_slug)); + const result = await kv.get(`${folderId}:dead`).catch(() => null) + return Boolean(result) + } + await migrateCollection({ db, mongooseCollection: Workspace, @@ -534,19 +569,19 @@ const main = async () => { // expired tokens can be removed // cannot use this uuid for the org id - return doc.environments.map((env, index) => { + return Promise.all(doc.environments.map(async (env, index) => { const id = uuidV4(); - envKv.put(env.slug, id); + await envKv.put(truncateAndSlugify(env.slug), id); return { id, name: env.name, - slug: env.slug, + slug: truncateAndSlugify(env.slug), position: index + 1, projectId: doc._id.toString(), createdAt: new Date(), updatedAt: new Date(), }; - }); + })); }, }); @@ -564,9 +599,9 @@ const main = async () => { doc.environments.map(async (env) => { const id = uuidV4(); - const envId = await getEnvId(doc._id.toString(), env.slug); + const envId = await getEnvId(doc._id.toString(), truncateAndSlugify(env.slug)); - const folderKv = getFolderKv(doc._id.toString(), env.slug); + const folderKv = getFolderKv(doc._id.toString(), truncateAndSlugify(env.slug)); await folderKv.put("root", id); return { id, @@ -589,9 +624,13 @@ const main = async () => { preProcessing: async (doc) => { // expired tokens can be removed // cannot use this uuid for the org id + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null); + if(!projectKvRes) return ; + const id = uuidV4(); - const senderId = (await userKv.get(doc.sender.toString()).catch(() => null)) || null; + const senderId = (await userKv.get(doc.sender.toString()).catch(() => null)); if (!senderId) return; const receiverId = await userKv.get(doc.receiver.toString()).catch(() => null); @@ -620,13 +659,17 @@ const main = async () => { returnKeys: ["id"], preProcessing: async (doc) => { const id = uuidV4(); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + await projectRoleKv.put(doc._id.toString(), id); return { id, name: doc.name, projectId: doc.workspace.toString(), description: doc.description, - slug: doc.slug, + slug: truncateAndSlugify(doc.slug), permissions: doc.permissions ? JSON.stringify(packRules(doc.permissions as any)) : null, @@ -645,6 +688,9 @@ const main = async () => { preProcessing: async (doc) => { const id = uuidV4(); + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + const userId = await userKv.get(doc.user.toString()).catch(() => null); if (!userId) return @@ -671,10 +717,16 @@ const main = async () => { postgresTableName: TableName.SecretFolder, returnKeys: ["id"], preProcessing: async (doc) => { - const folderKv = getFolderKv(doc.workspace.toString(), doc.environment); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const folderKv = getFolderKv(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + + const envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); const folders = flattenFolders(doc.nodes); + if (!folders) return + const pgFolder = await Promise.all( folders // already has been created @@ -708,8 +760,20 @@ const main = async () => { postgresTableName: TableName.SecretFolderVersion, returnKeys: ["id"], preProcessing: async (doc) => { - const folderKv = getFolderKv(doc.workspace.toString(), doc.environment); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const folderKv = getFolderKv(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + + // looping through env, we can come across envs that do not exist in the present state. + // This is because some folder snap shots were taken when that env slug exists but the same env slug was later deleted + let envId: string + try { + envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + }catch(e){ + return + } + const rootFolders = (doc?.nodes?.children || []).map( ({ name, version, id }) => ({ name, @@ -722,8 +786,18 @@ const main = async () => { rootFolders.map(async (folder) => { const { name, version } = folder; const id = uuidV4(); + await folderKv.put(getFolderVersionKey(folder.id, version), id); - const folderId = await folderKv.get(folder.id); + + // we are looking for each folder in folder versions and some we might not see + // because folderKv is ony keeps track of folders that are present in dashboard. So those we don't find, we'll add to same kv but add prefix `dead`. + // This way, we can handle dead ones we encounter in the future + const folderId = await folderKv.get(folder.id).catch(async () => { + const newFolderId = uuidV4(); + await folderKv.put(`${folder.id}:dead`, id) // we are adding dead folders that are not present in dashboard + return newFolderId + }); + return { name, version, @@ -746,14 +820,29 @@ const main = async () => { postgresTableName: TableName.SecretImport, returnKeys: ["id"], preProcessing: async (doc) => { + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + const envKv = envPKv.sublevel(doc.workspace.toString()); - const folderKv = getFolderKv(doc.workspace.toString(), doc.environment); - const folderId = await folderKv.get(doc.folderId); + + const folderKv = getFolderKv(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + + if (await checkIfFolderIsDangling(doc.workspace.toString(), truncateAndSlugify(doc.environment), doc.folderId)){ + return + } + + // case: when import is created for a given folder and later the folder is deleted AND the import with that folder ref is not deleted THEN this folder won't exist?? :( + const folderId = await folderKv.get(doc.folderId).catch(() => null); + + if (!folderId) return return Promise.all( doc.imports.map(async ({ environment, secretPath }, index) => { const id = uuidV4(); - const importEnv = await envKv.get(environment); + // case: when import is created but later the env slug is deleted HOWEVER the import with that env slug was not deleted :( + const importEnv = await envKv.get(truncateAndSlugify(environment)).catch(() => null); + if (!importEnv) return + return { id, folderId, @@ -764,7 +853,7 @@ const main = async () => { importEnv, importPath: secretPath, }; - }), + }).filter(Boolean) ); }, }); @@ -778,14 +867,22 @@ const main = async () => { preProcessing: async (doc) => { const id = uuidV4(); await tagKv.put(doc._id.toString(), id); + + // skip tags that have slugs that are empty + if (doc.slug.length < 1) { + return + } + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + const createdBy = await userKv.get(doc.user.toString()).catch(() => null); if (!createdBy) return; return { id, name: doc.name, - slug: doc.slug, + slug: truncateAndSlugify(doc.slug), color: doc.tagColor, projectId: doc.workspace.toString(), createdBy: createdBy || null, @@ -803,6 +900,9 @@ const main = async () => { preProcessing: async (doc) => { const id = uuidV4(); + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + return { id, projectId: doc.workspace.toString(), @@ -824,8 +924,29 @@ const main = async () => { postgresTableName: TableName.Secret, returnKeys: ["id"], preProcessing: async (doc) => { - const folderKv = getFolderKv(doc.workspace.toString(), doc.environment); - const folderId = await folderKv.get(doc.folder || "root"); + // when env slug is empty + if (doc.environment.length < 1) { + return + } + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + // case: we forgot to clean up secrets that belong to deleted env slugs + const isEnvFound = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)).catch(() => null) + if (!isEnvFound) return + + const folderKv = getFolderKv(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + + if (await checkIfFolderIsDangling(doc.workspace.toString(), truncateAndSlugify(doc.environment), doc.folder as string)){ + return + } + + // Case: if folder id doesn't exist and the root of the folder also doesn;t exist, THEN put the secret at the ROOT + const folderId = await folderKv.get(doc.folder || "root").catch(() => { + console.log("secret location unknown, moving secret to root of env_slug/project") + return "root" + }); const userId = doc.user ? await userKv.get(doc.user.toString()).catch(() => null) : null; if (!userId) return @@ -871,14 +992,33 @@ const main = async () => { return Promise.all( (doc.tags || [])?.map(async (tagId) => { const id = uuidV4(); - const secretId = await secKv.get(doc._id.toString()); - const secretTagId = await tagKv.get(tagId); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + // case: we forgot to clean up secrets that belong to deleted env slugs + const isEnvFound = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)).catch(() => null) + if (!isEnvFound) return + + if (await checkIfFolderIsDangling(doc.workspace.toString(), truncateAndSlugify(doc.environment), doc.folder as string)){ + return + } + + const userId = doc.user ? await userKv.get(doc.user.toString()).catch(() => null) : null; + if (!userId) return + + const secretId = await secKv.get(doc._id.toString()).catch((e) => { + throw e; + }); + const secretTagId = await tagKv.get(tagId).catch((e) => { + throw e; + }); return { id, [`${TableName.Secret}Id`]: secretId, [`${TableName.SecretTag}Id`]: secretTagId, - createdAt: new Date((doc as any).createdAt), - updatedAt: new Date((doc as any).updatedAt), + // createdAt: new Date((doc as any).createdAt), + // updatedAt: new Date((doc as any).updatedAt), }; }), ); @@ -892,17 +1032,44 @@ const main = async () => { postgresTableName: TableName.SecretVersion, returnKeys: ["id"], preProcessing: async (doc) => { - const folderKv = getFolderKv(doc.workspace.toString(), doc.environment); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); - const folderId = await folderKv.get(doc.folder || "root"); + // when env slug is empty + if (doc.environment.length < 1) { + return + } + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + // case: we forgot to clean up secrets that belong to deleted env slugs + const isEnvFound = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)).catch(() => null) + if (!isEnvFound) return + const folderKv = getFolderKv(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + + if (await checkIfFolderIsDangling(doc.workspace.toString(), truncateAndSlugify(doc.environment), doc.folder as string)){ + return + } + + // Case: if folder id doesn't exist and the root of the folder also doesn;t exist, THEN put the secret at the ROOT + const folderId = await folderKv.get(doc.folder || "root").catch(async () => { + console.log("secret location unknown, moving secret to root of env_slug/project") + return await folderKv.get("root") + }); + + const envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + const userId = doc.user ? await userKv.get(doc.user.toString()).catch(() => null) : null; if (!userId) return const id = uuidV4(); await secVerKv.put(doc._id.toString(), id); - const secretId = await secKv.get(doc.secret.toString()); + const secretId = await secKv.get(doc.secret.toString()).catch(async (e) =>{ + const newId = uuidV4(); + await secKv.put(`${doc.secret.toString()}:dead`, newId) + return newId; + }); + // comment and reminder are not saved in secret version of mongo return { id, @@ -935,10 +1102,20 @@ const main = async () => { returnKeys: ["id"], preProcessing: async (doc) => { const id = uuidV4(); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + await projectBotKv.put(doc._id.toString(), id); - const botKey = await BotKey.findOne({ + + // case: skip bots that are inactive (skipped specifically because 6388653a200193a667c7e3f3 has two records, one active one not) + let botKey = await BotKey.findOne({ workspace: doc.workspace, - }).lean(); + $or: [ + { isActive: true }, + { isActive: false } + ] + }).sort({ isActive: -1 }).lean(); const senderId = botKey?.sender ? await userKv.get(botKey.sender.toString()).catch(() => null) @@ -975,6 +1152,10 @@ const main = async () => { preProcessing: async (doc) => { const id = uuidV4(); await integrationAuthKv.put(doc._id.toString(), id); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + return { id, projectId: doc.workspace.toString(), @@ -1009,7 +1190,11 @@ const main = async () => { returnKeys: ["id"], preProcessing: async (doc) => { const id = uuidV4(); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); const integrationAuthId = await integrationAuthKv.get( doc.integrationAuth.toString(), ); @@ -1047,6 +1232,9 @@ const main = async () => { preProcessing: async (doc) => { const id = uuidV4(); + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + const userId = await userKv.get(doc.user.toString()).catch(() => null); if (!userId) return; @@ -1076,7 +1264,11 @@ const main = async () => { returnKeys: ["id"], preProcessing: async (doc) => { const id = uuidV4(); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); return { id, iv: doc.iv, @@ -1240,6 +1432,10 @@ const main = async () => { returnKeys: ["id"], preProcessing: async (doc) => { const id = uuidV4(); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + const identityId = await identityKv.get(doc.identity.toString()); const roleId = doc.customRole ? await projectRoleKv.get(doc.customRole.toString()) @@ -1265,8 +1461,12 @@ const main = async () => { returnKeys: ["id"], preProcessing: async (doc) => { const id = uuidV4(); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); - sapKv.put(doc._id.toString(), id); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + await sapKv.put(doc._id.toString(), id); return { id, @@ -1316,7 +1516,11 @@ const main = async () => { preProcessing: async (doc) => { const id = uuidV4(); await secRotationKv.put(doc._id.toString(), id); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); return { id, @@ -1535,6 +1739,9 @@ const main = async () => { preProcessing: async (doc) => { const id = uuidV4(); + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + return { id, projectId: doc.workspace.toString(), @@ -1558,8 +1765,12 @@ const main = async () => { preProcessing: async (doc) => { const id = uuidV4(); await snapKv.put(doc._id.toString(), id); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); - const folderKv = getFolderKv(doc.workspace.toString(), doc.environment); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + const folderKv = getFolderKv(doc.workspace.toString(), truncateAndSlugify(doc.environment)); const folderId = await folderKv.get(doc.folderId).catch(async () => { // this folder may not exist now in tree then create a new id and assign it const newId = uuidV4(); @@ -1585,7 +1796,11 @@ const main = async () => { returnKeys: ["id"], preProcessing: async (doc) => { const snapshotId = await snapKv.get(doc._id.toString()); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); return Promise.all( doc.secretVersions.map(async (secVer) => { @@ -1612,8 +1827,12 @@ const main = async () => { returnKeys: ["id"], preProcessing: async (doc) => { const snapshotId = await snapKv.get(doc._id.toString()); - const envId = await getEnvId(doc.workspace.toString(), doc.environment); - const folderKv = getFolderKv(doc.workspace.toString(), doc.environment); + + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const envId = await getEnvId(doc.workspace.toString(), truncateAndSlugify(doc.environment)); + const folderKv = getFolderKv(doc.workspace.toString(), truncateAndSlugify(doc.environment)); const folderVersion = await FolderVersion.findById(doc.folderVersion); if (!folderVersion) return; @@ -1647,7 +1866,10 @@ const main = async () => { await sarKv.put(doc._id.toString(), id); const policyId = await sapKv.get(doc.policy.toString()); - const folderKv = getFolderKv(doc.workspace.toString(), doc.environment); + const projectKvRes = await projectKv.get(doc.workspace.toString()).catch(() => null) + if (!projectKvRes) return + + const folderKv = getFolderKv(doc.workspace.toString(), truncateAndSlugify(doc.environment)); const folderId = await folderKv.get(doc.folderId); const committerId = await projectMembKv.get(doc.committer.toString()); @@ -1660,7 +1882,7 @@ const main = async () => { hasMerged: doc.hasMerged, status: doc.status, conflicts: JSON.stringify(doc.conflicts), - slug: doc.slug, + slug: truncateAndSlugify(doc.slug), folderId, committerId, statusChangeBy,