PIT: Add initialization and checkpoint logic

This commit is contained in:
carlosmonastyrski
2025-05-08 09:41:01 -03:00
parent e58dbe853e
commit e8d424bbb0
7 changed files with 260 additions and 24 deletions

View File

@@ -0,0 +1,26 @@
import { Knex } from "knex";
import { ProjectType, TableName } from "../schemas";
import { getMigrationPITServices } from "./utils/services";
export async function up(knex: Knex): Promise<void> {
const hasFolderCommitTable = await knex.schema.hasTable(TableName.FolderCommit);
if (hasFolderCommitTable) {
const { folderCommitService } = await getMigrationPITServices({ db: knex });
const projects = await knex(TableName.Project).where({ version: 3, type: ProjectType.SecretManager }).select("id");
await knex.transaction(async (tx) => {
for (const project of projects) {
// eslint-disable-next-line no-await-in-loop
await folderCommitService.initializeProject(project.id, tx);
}
});
}
}
export async function down(knex: Knex): Promise<void> {
const hasFolderCommitTable = await knex.schema.hasTable(TableName.FolderCommit);
if (hasFolderCommitTable) {
// delete all existing entries
await knex(TableName.FolderCommit).del();
}
}

View File

@@ -3,12 +3,23 @@ import { Knex } from "knex";
import { initializeHsmModule } from "@app/ee/services/hsm/hsm-fns";
import { hsmServiceFactory } from "@app/ee/services/hsm/hsm-service";
import { TKeyStoreFactory } from "@app/keystore/keystore";
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 { 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 { 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";
import { kmsRootConfigDALFactory } from "@app/services/kms/kms-root-config-dal";
import { kmsServiceFactory } from "@app/services/kms/kms-service";
import { orgDALFactory } from "@app/services/org/org-dal";
import { projectDALFactory } from "@app/services/project/project-dal";
import { secretFolderDALFactory } from "@app/services/secret-folder/secret-folder-dal";
import { secretFolderVersionDALFactory } from "@app/services/secret-folder/secret-folder-version-dal";
import { secretVersionV2BridgeDALFactory } from "@app/services/secret-v2-bridge/secret-version-dal";
import { userDALFactory } from "@app/services/user/user-dal";
import { TMigrationEnvConfig } from "./env-config";
@@ -50,3 +61,33 @@ export const getMigrationEncryptionServices = async ({ envConfig, db, keyStore }
return { kmsService };
};
export const getMigrationPITServices = async ({ db }: { db: Knex }) => {
const projectDAL = projectDALFactory(db);
const folderCommitDAL = folderCommitDALFactory(db);
const folderCommitChangesDAL = folderCommitChangesDALFactory(db);
const folderCheckpointDAL = folderCheckpointDALFactory(db);
const folderTreeCheckpointDAL = folderTreeCheckpointDALFactory(db);
const userDAL = userDALFactory(db);
const identityDAL = identityDALFactory(db);
const folderDAL = secretFolderDALFactory(db);
const folderVersionDAL = secretFolderVersionDALFactory(db);
const secretVersionV2BridgeDAL = secretVersionV2BridgeDALFactory(db);
const folderCheckpointResourcesDAL = folderCheckpointResourcesDALFactory(db);
const folderCommitService = folderCommitServiceFactory({
folderCommitDAL,
folderCommitChangesDAL,
folderCheckpointDAL,
folderTreeCheckpointDAL,
userDAL,
identityDAL,
folderDAL,
folderVersionDAL,
secretVersionV2BridgeDAL,
projectDAL,
folderCheckpointResourcesDAL
});
return { folderCommitService };
};

View File

@@ -228,6 +228,9 @@ const envSchema = z
DATADOG_SERVICE: zpStr(z.string().optional().default("infisical-core")),
DATADOG_HOSTNAME: zpStr(z.string().optional()),
// PIT
CHECKPOINT_WINDOW: zpStr(z.string().optional().default("10")),
/* CORS ----------------------------------------------------------------------------- */
CORS_ALLOWED_ORIGINS: zpStr(

View File

@@ -141,6 +141,7 @@ import { externalGroupOrgRoleMappingServiceFactory } from "@app/services/externa
import { externalMigrationQueueFactory } from "@app/services/external-migration/external-migration-queue";
import { externalMigrationServiceFactory } from "@app/services/external-migration/external-migration-service";
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 { folderCommitServiceFactory } from "@app/services/folder-commit/folder-commit-service";
import { folderCommitChangesDALFactory } from "@app/services/folder-commit-changes/folder-commit-changes-dal";
@@ -559,6 +560,7 @@ export const registerRoutes = async (
const folderCommitChangesDAL = folderCommitChangesDALFactory(db);
const folderCheckpointDAL = folderCheckpointDALFactory(db);
const folderCheckpointResourcesDAL = folderCheckpointResourcesDALFactory(db);
const folderTreeCheckpointDAL = folderTreeCheckpointDALFactory(db);
const folderCommitDAL = folderCommitDALFactory(db);
const folderCommitService = folderCommitServiceFactory({
@@ -567,7 +569,12 @@ export const registerRoutes = async (
folderCheckpointDAL,
folderTreeCheckpointDAL,
userDAL,
identityDAL
identityDAL,
folderDAL,
folderVersionDAL,
secretVersionV2BridgeDAL,
projectDAL,
folderCheckpointResourcesDAL
});
const scimService = scimServiceFactory({
licenseService,

View File

@@ -23,24 +23,12 @@ export const folderCommitDALFactory = (db: TDbClient) => {
}
};
const findById = async (id: string, tx?: Knex): Promise<TFolderCommits | undefined> => {
try {
const doc = await (tx || db.replicaNode())(TableName.FolderCommit)
.where({ id })
.select(selectAllTableCols(TableName.FolderCommit))
.first();
return doc;
} catch (error) {
throw new DatabaseError({ error, name: "FindById" });
}
};
const findLatestCommit = async (folderId: string, tx?: Knex): Promise<TFolderCommits | undefined> => {
try {
const doc = await (tx || db.replicaNode())(TableName.FolderCommit)
.where({ folderId })
.select(selectAllTableCols(TableName.FolderCommit))
.orderBy("createdAt", "desc")
.orderBy("commitId", "desc")
.first();
return doc;
} catch (error) {
@@ -48,10 +36,30 @@ export const folderCommitDALFactory = (db: TDbClient) => {
}
};
const getNumberOfCommitsSince = async (folderId: string, folderCommitId: string, tx?: Knex): Promise<number> => {
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({ folderId })
.where("commitId", ">", referencedCommit.commitId)
.count();
return Number(doc?.[0].count);
}
return 0;
} catch (error) {
throw new DatabaseError({ error, name: "getNumberOfCommitsSince" });
}
};
return {
...restOfOrm,
findByFolderId,
findById,
findLatestCommit
findLatestCommit,
getNumberOfCommitsSince
};
};

View File

@@ -1,28 +1,37 @@
import { Knex } from "knex";
import { TSecretFolders } from "@app/db/schemas";
import { getConfig } from "@app/lib/config/env";
import { BadRequestError, DatabaseError, NotFoundError } from "@app/lib/errors";
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 { 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 { TSecretVersionV2DALFactory } from "../secret-v2-bridge/secret-version-dal";
import { TUserDALFactory } from "../user/user-dal";
import { TFolderCommitDALFactory } from "./folder-commit-dal";
type TFolderCommitServiceFactoryDep = {
folderCommitDAL: Pick<
TFolderCommitDALFactory,
"create" | "findById" | "findByFolderId" | "findLatestCommit" | "transaction"
"create" | "findById" | "findByFolderId" | "findLatestCommit" | "transaction" | "getNumberOfCommitsSince"
>;
folderCommitChangesDAL: Pick<TFolderCommitChangesDALFactory, "create" | "findByCommitId" | "insertMany">;
folderCheckpointDAL: Pick<TFolderCheckpointDALFactory, "create" | "findByFolderId" | "findLatestByFolderId">;
folderTreeCheckpointDAL: Pick<
TFolderTreeCheckpointDALFactory,
"create" | "findByProjectId" | "findLatestByProjectId"
>;
folderCheckpointResourcesDAL: Pick<TFolderCheckpointResourcesDALFactory, "insertMany">;
folderTreeCheckpointDAL: Pick<TFolderTreeCheckpointDALFactory, "findByProjectId" | "findLatestByProjectId">;
userDAL: Pick<TUserDALFactory, "findById">;
identityDAL: Pick<TIdentityDALFactory, "findById">;
folderDAL: Pick<TSecretFolderDALFactory, "findByParentId" | "findByProjectId">;
folderVersionDAL: Pick<TSecretFolderVersionDALFactory, "findLatestFolderVersions">;
secretVersionV2BridgeDAL: Pick<TSecretVersionV2DALFactory, "findLatestVersionByFolderId">;
projectDAL: Pick<TProjectDALFactory, "findById">;
};
export type TCreateCommitDTO = {
@@ -54,9 +63,72 @@ export const folderCommitServiceFactory = ({
folderCommitChangesDAL,
folderCheckpointDAL,
folderTreeCheckpointDAL,
folderCheckpointResourcesDAL,
userDAL,
identityDAL
identityDAL,
folderDAL,
folderVersionDAL,
secretVersionV2BridgeDAL,
projectDAL
}: TFolderCommitServiceFactoryDep) => {
const appCfg = getConfig();
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;
};
const createFolderCheckpoint = async ({
folderId,
folderCommitId,
force = false,
tx
}: {
folderId: string;
folderCommitId?: string;
force?: boolean;
tx?: Knex;
}) => {
let latestCommitId = folderCommitId;
if (!latestCommitId) {
latestCommitId = (await folderCheckpointDAL.findLatestByFolderId(folderId, tx))?.folderCommitId;
}
if (!latestCommitId) {
throw new BadRequestError({ message: "Latest commit ID not found" });
return;
}
if (!force) {
const commitsSinceLastCheckpoint = await folderCommitDAL.getNumberOfCommitsSince(folderId, latestCommitId, tx);
if (commitsSinceLastCheckpoint < Number(appCfg.CHECKPOINT_WINDOW)) {
return;
}
}
const checkpointResources = await getFolderResources(folderId, tx);
if (checkpointResources.length > 0) {
const newCheckpoint = await folderCheckpointDAL.create(
{
folderCommitId: latestCommitId
},
tx
);
await folderCheckpointResourcesDAL.insertMany(
checkpointResources.map((resource) => ({ folderCheckpointId: newCheckpoint.id, ...resource })),
tx
);
}
};
const createCommit = async (data: TCreateCommitDTO, tx?: Knex) => {
const metadata = data.actor.metadata || {};
try {
@@ -87,6 +159,7 @@ export const folderCommitServiceFactory = ({
tx
);
await createFolderCheckpoint({ folderId: data.folderId, folderCommitId: newCommit.id, tx });
return newCommit;
} catch (error) {
throw new DatabaseError({ error, name: "CreateCommit" });
@@ -149,6 +222,69 @@ export const folderCommitServiceFactory = ({
return folderTreeCheckpointDAL.findLatestByProjectId(projectId, tx);
};
const initializeFolder = async (folderId: string, tx?: Knex) => {
const folderResources = await getFolderResources(folderId, tx);
const changes = folderResources.map((resource) => ({ type: "add", ...resource }));
if (changes.length > 0) {
const newCommit = await createCommit(
{
actor: {
type: ActorType.PLATFORM
},
message: "Initialized folder",
folderId,
changes
},
tx
);
await createFolderCheckpoint({ folderId, folderCommitId: newCommit.id, force: true, tx });
}
};
function sortFoldersByHierarchy(folders: TSecretFolders[]) {
// Create a map for quick lookup of children by parent ID
const childrenMap: Map<string | null, TSecretFolders[]> = new Map();
folders.forEach((folder) => {
const { parentId } = folder;
if (!childrenMap.has(parentId || null)) {
childrenMap.set(parentId || null, []);
}
childrenMap.get(parentId || null)?.push(folder);
});
// Start with root folders (null parentId)
const result = [];
const rootFolders = childrenMap.get(null) || [];
// Process each level of the hierarchy
let currentLevel = rootFolders;
result.push(...currentLevel);
while (currentLevel.length > 0) {
const nextLevel = [];
for (const folder of currentLevel) {
const children = childrenMap.get(folder.id) || [];
nextLevel.push(...children);
}
result.push(...nextLevel);
currentLevel = nextLevel;
}
return result;
}
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)));
};
return {
createCommit,
addCommitChange,
@@ -158,7 +294,10 @@ export const folderCommitServiceFactory = ({
getCheckpointsByFolderId,
getLatestCheckpoint,
getTreeCheckpointsByProjectId,
getLatestTreeCheckpoint
getLatestTreeCheckpoint,
initializeFolder,
initializeProject,
createFolderCheckpoint
};
};

View File

@@ -488,6 +488,17 @@ export const secretFolderDALFactory = (db: TDbClient) => {
}
};
const findByParentId = async (parentId: string, tx?: Knex) => {
try {
const folders = await (tx || db.replicaNode())(TableName.SecretFolder)
.where({ parentId })
.select(selectAllTableCols(TableName.SecretFolder));
return folders;
} catch (error) {
throw new DatabaseError({ error, name: "findByParentId" });
}
};
return {
...secretFolderOrm,
update,
@@ -499,6 +510,7 @@ export const secretFolderDALFactory = (db: TDbClient) => {
findClosestFolder,
findByProjectId,
findByMultiEnv,
findByEnvsDeep
findByEnvsDeep,
findByParentId
};
};