Update migration

This commit is contained in:
x032205
2025-09-05 13:16:18 -04:00
parent 17b6ab0db0
commit 3a75a50d0b
@@ -24,96 +24,98 @@ export async function up(knex: Knex): Promise<void> {
t.string("url").nullable().alter(); t.string("url").nullable().alter();
}); });
const superAdminDAL = superAdminDALFactory(knex); if (!hasEncryptedCredentials) {
const envConfig = await getMigrationEnvConfig(superAdminDAL); const superAdminDAL = superAdminDALFactory(knex);
const keyStore = inMemoryKeyStore(); const envConfig = await getMigrationEnvConfig(superAdminDAL);
const keyStore = inMemoryKeyStore();
const { kmsService } = await getMigrationEncryptionServices({ envConfig, keyStore, db: knex }); const { kmsService } = await getMigrationEncryptionServices({ envConfig, keyStore, db: knex });
const orgEncryptionRingBuffer = const orgEncryptionRingBuffer =
createCircularCache<Awaited<ReturnType<(typeof kmsService)["createCipherPairWithDataKey"]>>>(25); createCircularCache<Awaited<ReturnType<(typeof kmsService)["createCipherPairWithDataKey"]>>>(25);
const logStreams = await knex(TableName.AuditLogStream).select( const logStreams = await knex(TableName.AuditLogStream).select(
"id", "id",
"orgId", "orgId",
"url", "url",
"encryptedHeadersAlgorithm", "encryptedHeadersAlgorithm",
"encryptedHeadersCiphertext", "encryptedHeadersCiphertext",
"encryptedHeadersIV", "encryptedHeadersIV",
"encryptedHeadersKeyEncoding", "encryptedHeadersKeyEncoding",
"encryptedHeadersTag" "encryptedHeadersTag"
); );
const updatedLogStreams = await Promise.all( const updatedLogStreams = await Promise.all(
logStreams.map(async (el) => { logStreams.map(async (el) => {
let orgKmsService = orgEncryptionRingBuffer.getItem(el.orgId); let orgKmsService = orgEncryptionRingBuffer.getItem(el.orgId);
if (!orgKmsService) { if (!orgKmsService) {
orgKmsService = await kmsService.createCipherPairWithDataKey( orgKmsService = await kmsService.createCipherPairWithDataKey(
{ {
type: KmsDataKey.Organization, type: KmsDataKey.Organization,
orgId: el.orgId orgId: el.orgId
}, },
knex knex
); );
orgEncryptionRingBuffer.push(el.orgId, orgKmsService); orgEncryptionRingBuffer.push(el.orgId, orgKmsService);
} }
const provider = "custom"; const provider = "custom";
let credentials; let credentials;
if ( if (
el.encryptedHeadersTag && el.encryptedHeadersTag &&
el.encryptedHeadersIV && el.encryptedHeadersIV &&
el.encryptedHeadersCiphertext && el.encryptedHeadersCiphertext &&
el.encryptedHeadersKeyEncoding el.encryptedHeadersKeyEncoding
) { ) {
const decryptedHeaders = crypto const decryptedHeaders = crypto
.encryption() .encryption()
.symmetric() .symmetric()
.decryptWithRootEncryptionKey({ .decryptWithRootEncryptionKey({
tag: el.encryptedHeadersTag, tag: el.encryptedHeadersTag,
iv: el.encryptedHeadersIV, iv: el.encryptedHeadersIV,
ciphertext: el.encryptedHeadersCiphertext, ciphertext: el.encryptedHeadersCiphertext,
keyEncoding: el.encryptedHeadersKeyEncoding as SecretKeyEncoding keyEncoding: el.encryptedHeadersKeyEncoding as SecretKeyEncoding
}); });
credentials = { credentials = {
url: el.url,
headers: JSON.parse(decryptedHeaders)
};
} else {
credentials = {
url: el.url,
headers: []
};
}
const encryptedCredentials = orgKmsService.encryptor({
plainText: Buffer.from(JSON.stringify(credentials), "utf8")
}).cipherTextBlob;
return {
id: el.id,
orgId: el.orgId,
url: el.url, url: el.url,
headers: JSON.parse(decryptedHeaders) provider,
encryptedCredentials
}; };
} else { })
credentials = { );
url: el.url,
headers: []
};
}
const encryptedCredentials = orgKmsService.encryptor({ for (let i = 0; i < updatedLogStreams.length; i += BATCH_SIZE) {
plainText: Buffer.from(JSON.stringify(credentials), "utf8") // eslint-disable-next-line no-await-in-loop
}).cipherTextBlob; await knex(TableName.AuditLogStream)
.insert(updatedLogStreams.slice(i, i + BATCH_SIZE))
.onConflict("id")
.merge();
}
return { await knex.schema.alterTable(TableName.AuditLogStream, (t) => {
id: el.id, t.binary("encryptedCredentials").notNullable().alter();
orgId: el.orgId, });
url: el.url,
provider,
encryptedCredentials
};
})
);
for (let i = 0; i < updatedLogStreams.length; i += BATCH_SIZE) {
// eslint-disable-next-line no-await-in-loop
await knex(TableName.AuditLogStream)
.insert(updatedLogStreams.slice(i, i + BATCH_SIZE))
.onConflict("id")
.merge();
} }
await knex.schema.alterTable(TableName.AuditLogStream, (t) => {
t.binary("encryptedCredentials").notNullable().alter();
});
} }
} }
@@ -125,90 +127,95 @@ export async function up(knex: Knex): Promise<void> {
export async function down(knex: Knex): Promise<void> { export async function down(knex: Knex): Promise<void> {
if (await knex.schema.hasTable(TableName.AuditLogStream)) { if (await knex.schema.hasTable(TableName.AuditLogStream)) {
const superAdminDAL = superAdminDALFactory(knex); const hasProvider = await knex.schema.hasColumn(TableName.AuditLogStream, "provider");
const envConfig = await getMigrationEnvConfig(superAdminDAL); const hasEncryptedCredentials = await knex.schema.hasColumn(TableName.AuditLogStream, "encryptedCredentials");
const keyStore = inMemoryKeyStore();
const { kmsService } = await getMigrationEncryptionServices({ envConfig, keyStore, db: knex }); if (hasEncryptedCredentials) {
const superAdminDAL = superAdminDALFactory(knex);
const envConfig = await getMigrationEnvConfig(superAdminDAL);
const keyStore = inMemoryKeyStore();
const orgEncryptionRingBuffer = const { kmsService } = await getMigrationEncryptionServices({ envConfig, keyStore, db: knex });
createCircularCache<Awaited<ReturnType<(typeof kmsService)["createCipherPairWithDataKey"]>>>(25);
const logStreamsToRevert = await knex(TableName.AuditLogStream) const orgEncryptionRingBuffer =
.select("id", "orgId", "encryptedCredentials") createCircularCache<Awaited<ReturnType<(typeof kmsService)["createCipherPairWithDataKey"]>>>(25);
.where("provider", "custom")
.whereNotNull("encryptedCredentials");
const updatedLogStreams = await Promise.all( const logStreamsToRevert = await knex(TableName.AuditLogStream)
logStreamsToRevert.map(async (el) => { .select("id", "orgId", "encryptedCredentials")
let orgKmsService = orgEncryptionRingBuffer.getItem(el.orgId); .where("provider", "custom")
if (!orgKmsService) { .whereNotNull("encryptedCredentials");
orgKmsService = await kmsService.createCipherPairWithDataKey(
{
type: KmsDataKey.Organization,
orgId: el.orgId
},
knex
);
orgEncryptionRingBuffer.push(el.orgId, orgKmsService);
}
const decryptedCredentials = orgKmsService const updatedLogStreams = await Promise.all(
.decryptor({ logStreamsToRevert.map(async (el) => {
cipherTextBlob: el.encryptedCredentials let orgKmsService = orgEncryptionRingBuffer.getItem(el.orgId);
}) if (!orgKmsService) {
.toString(); orgKmsService = await kmsService.createCipherPairWithDataKey(
{
type: KmsDataKey.Organization,
orgId: el.orgId
},
knex
);
orgEncryptionRingBuffer.push(el.orgId, orgKmsService);
}
const credentials: { url: string; headers: { key: string; value: string }[] } = const decryptedCredentials = orgKmsService
JSON.parse(decryptedCredentials); .decryptor({
cipherTextBlob: el.encryptedCredentials
})
.toString();
const originalUrl: string = credentials.url; const credentials: { url: string; headers: { key: string; value: string }[] } =
JSON.parse(decryptedCredentials);
const encryptedHeadersResult = crypto const originalUrl: string = credentials.url;
.encryption()
.symmetric()
.encryptWithRootEncryptionKey(JSON.stringify(credentials.headers), envConfig);
const encryptedHeadersAlgorithm: string = encryptedHeadersResult.algorithm; const encryptedHeadersResult = crypto
const encryptedHeadersCiphertext: string = encryptedHeadersResult.ciphertext; .encryption()
const encryptedHeadersIV: string = encryptedHeadersResult.iv; .symmetric()
const encryptedHeadersKeyEncoding: string = encryptedHeadersResult.encoding; .encryptWithRootEncryptionKey(JSON.stringify(credentials.headers), envConfig);
const encryptedHeadersTag: string = encryptedHeadersResult.tag;
return { const encryptedHeadersAlgorithm: string = encryptedHeadersResult.algorithm;
id: el.id, const encryptedHeadersCiphertext: string = encryptedHeadersResult.ciphertext;
orgId: el.orgId, const encryptedHeadersIV: string = encryptedHeadersResult.iv;
encryptedCredentials: el.encryptedCredentials, const encryptedHeadersKeyEncoding: string = encryptedHeadersResult.encoding;
const encryptedHeadersTag: string = encryptedHeadersResult.tag;
url: originalUrl, return {
encryptedHeadersAlgorithm, id: el.id,
encryptedHeadersCiphertext, orgId: el.orgId,
encryptedHeadersIV, encryptedCredentials: el.encryptedCredentials,
encryptedHeadersKeyEncoding,
encryptedHeadersTag url: originalUrl,
}; encryptedHeadersAlgorithm,
}) encryptedHeadersCiphertext,
); encryptedHeadersIV,
encryptedHeadersKeyEncoding,
encryptedHeadersTag
};
})
);
for (let i = 0; i < updatedLogStreams.length; i += BATCH_SIZE) {
// eslint-disable-next-line no-await-in-loop
await knex(TableName.AuditLogStream)
.insert(updatedLogStreams.slice(i, i + BATCH_SIZE))
.onConflict("id")
.merge();
}
for (let i = 0; i < updatedLogStreams.length; i += BATCH_SIZE) {
// eslint-disable-next-line no-await-in-loop
await knex(TableName.AuditLogStream) await knex(TableName.AuditLogStream)
.insert(updatedLogStreams.slice(i, i + BATCH_SIZE)) .where((qb) => {
.onConflict("id") void qb.whereNot("provider", "custom").orWhereNull("url");
.merge(); })
.del();
} }
await knex(TableName.AuditLogStream)
.where((qb) => {
void qb.whereNot("provider", "custom").orWhereNull("url");
})
.del();
await knex.schema.alterTable(TableName.AuditLogStream, (t) => { await knex.schema.alterTable(TableName.AuditLogStream, (t) => {
t.string("url").notNullable().alter(); t.string("url").notNullable().alter();
t.dropColumn("provider"); if (hasProvider) t.dropColumn("provider");
t.dropColumn("encryptedCredentials"); if (hasEncryptedCredentials) t.dropColumn("encryptedCredentials");
}); });
} }
} }