Execute reviewed tenant merges safely

This commit is contained in:
2026-08-06 12:34:52 +02:00
parent a2c5ad5dd9
commit e2fdfe9c0a
5 changed files with 440 additions and 11 deletions

View File

@@ -1,9 +1,12 @@
import { pool } from "../../db"
import type { TenantFullExport } from "./tenantFullExport"
import { buildTenantMergePlan, type TenantMergePlan, type TenantMergeTableMetadata } from "./tenantMergePlan"
import { prepareTenantFullExportRowsForImport, type TenantFullExport } from "./tenantFullExport"
import { buildTenantMergePlan, type TenantMergeDecision, type TenantMergePlan, type TenantMergeTableMetadata } from "./tenantMergePlan"
type MergeDatabaseMetadata = TenantMergeTableMetadata & {
columns: string[]
jsonColumns: Set<string>
generatedColumns: Set<string>
foreignKeys: { column: string, referencedTable: string, referencedColumn: string }[]
}
const quoteIdent = (value: string) => `"${value.replace(/"/g, '""')}"`
@@ -11,11 +14,23 @@ const globalNaturalKeyTables = new Set(["accounts", "units", "countrys", "citys"
const loadMergeMetadata = async (client: any) => {
const columnsResult = await client.query(`
select table_name, column_name
select table_name, column_name, data_type, is_generated
from information_schema.columns
where table_schema = 'public'
order by table_name, ordinal_position
`)
const foreignKeyResult = await client.query(`
select tc.table_name, kcu.column_name, ccu.table_name as referenced_table,
ccu.column_name as referenced_column
from information_schema.table_constraints tc
join information_schema.key_column_usage kcu
on tc.constraint_name = kcu.constraint_name
and tc.constraint_schema = kcu.constraint_schema
join information_schema.constraint_column_usage ccu
on tc.constraint_name = ccu.constraint_name
and tc.constraint_schema = ccu.constraint_schema
where tc.table_schema = 'public' and tc.constraint_type = 'FOREIGN KEY'
`)
const primaryKeyResult = await client.query(`
select tc.table_name, kcu.column_name
from information_schema.table_constraints tc
@@ -28,13 +43,23 @@ const loadMergeMetadata = async (client: any) => {
const metadata: Record<string, MergeDatabaseMetadata> = {}
for (const row of columnsResult.rows) {
metadata[row.table_name] ||= { columns: [], primaryKey: [] }
metadata[row.table_name] ||= { columns: [], primaryKey: [], jsonColumns: new Set(), generatedColumns: new Set(), foreignKeys: [] }
metadata[row.table_name].columns.push(row.column_name)
if (row.data_type === "json" || row.data_type === "jsonb") metadata[row.table_name].jsonColumns.add(row.column_name)
if (row.is_generated === "ALWAYS") metadata[row.table_name].generatedColumns.add(row.column_name)
}
for (const row of primaryKeyResult.rows) {
metadata[row.table_name] ||= { columns: [], primaryKey: [] }
metadata[row.table_name] ||= { columns: [], primaryKey: [], jsonColumns: new Set(), generatedColumns: new Set(), foreignKeys: [] }
metadata[row.table_name].primaryKey.push(row.column_name)
}
for (const row of foreignKeyResult.rows) {
metadata[row.table_name] ||= { columns: [], primaryKey: [], jsonColumns: new Set(), generatedColumns: new Set(), foreignKeys: [] }
metadata[row.table_name].foreignKeys.push({
column: row.column_name,
referencedTable: row.referenced_table,
referencedColumn: row.referenced_column,
})
}
return metadata
}
@@ -122,3 +147,225 @@ export const createTenantMergeDryRun = async (
client.release()
}
}
const sensitiveColumns = new Set([
"iban_encrypted",
"bic_encrypted",
"bank_name_encrypted",
"__plainIban",
"__plainBic",
"__plainBankName",
"password_hash",
"passwordHash",
"token_hash",
"tokenHash",
])
const redactRow = (row: Record<string, any> | null) => row && Object.fromEntries(
Object.entries(row).map(([column, value]) => [column, sensitiveColumns.has(column) && value ? "***" : value])
)
export const sanitizeTenantMergePlan = (plan: TenantMergePlan): TenantMergePlan => ({
...plan,
items: plan.items.map((item) => ({
...item,
sourceRow: redactRow(item.sourceRow) || {},
targetRow: redactRow(item.targetRow),
})),
})
const prepareValue = (value: any, isJson: boolean) => {
if (!isJson || value === null || typeof value === "undefined" || typeof value === "string") return value
return JSON.stringify(value)
}
const rowIdentity = (row: Record<string, any>, columns: string[]) =>
columns.map((column) => row[column]).join("\0")
const topologicalTableOrder = (tables: string[], metadata: Record<string, MergeDatabaseMetadata>) => {
const remaining = new Set(tables)
const ordered: string[] = []
while (remaining.size) {
const ready = Array.from(remaining).filter((table) =>
metadata[table].foreignKeys.every((foreignKey) => !remaining.has(foreignKey.referencedTable) || foreignKey.referencedTable === table)
)
const next = ready.length ? ready.sort() : [Array.from(remaining).sort()[0]]
for (const table of next) {
remaining.delete(table)
ordered.push(table)
}
}
return ordered
}
export type TenantMergeExecutionResult = {
imported: number
updated: number
retained: number
remappedIds: number
tables: Record<string, { imported: number, updated: number, retained: number }>
}
export const executeTenantMergeWithClient = async (
client: any,
rawExportData: TenantFullExport,
targetTenantId: number,
decisions: Record<string, TenantMergeDecision>
): Promise<TenantMergeExecutionResult> => {
const metadata = await loadMergeMetadata(client)
const sourceTables = remapSourceTenant(rawExportData, targetTenantId)
const preparedExport: TenantFullExport = prepareTenantFullExportRowsForImport({
...rawExportData,
tenantId: targetTenantId,
tables: Object.fromEntries(Object.entries(sourceTables).map(([table, rows]) => [table, rows.map((row) => ({ ...row }))])),
})
const targetTables: Record<string, Record<string, any>[]> = {}
for (const [table, rows] of Object.entries(preparedExport.tables)) {
if (metadata[table]) targetTables[table] = await loadTargetRows(client, table, rows, metadata[table], targetTenantId)
}
const plan = buildTenantMergePlan(preparedExport.tables, targetTables, metadata)
const idMaps = new Map<string, Map<any, any>>()
const selected = plan.items.filter((item) => {
const decision = decisions[item.id] || item.defaultDecision
return decision === "source" && item.kind !== "existing"
})
for (const item of plan.items) {
const tableMetadata = metadata[item.table]
if (tableMetadata?.primaryKey.length !== 1 || !item.targetRow) continue
const key = tableMetadata.primaryKey[0]
const sourceId = item.sourceRow[key]
const targetId = item.targetRow[key]
if (sourceId !== null && typeof sourceId !== "undefined" && targetId !== null && typeof targetId !== "undefined") {
if (!idMaps.has(item.table)) idMaps.set(item.table, new Map())
idMaps.get(item.table)!.set(sourceId, targetId)
}
}
await client.query("begin")
await client.query("set local session_replication_role = replica")
try {
for (const item of selected.filter((entry) => entry.kind === "id_collision")) {
const tableMetadata = metadata[item.table]
if (tableMetadata.primaryKey.length !== 1) throw new Error(`ID-Kollision in ${item.table} kann nicht automatisch aufgelöst werden`)
const key = tableMetadata.primaryKey[0]
const sequenceResult = await client.query("select pg_get_serial_sequence($1, $2) as sequence_name", [`public.${item.table}`, key])
const sequenceName = sequenceResult.rows[0]?.sequence_name
if (!sequenceName) throw new Error(`Keine Sequenz für ID-Kollision in ${item.table}.${key} gefunden`)
const allocated = await client.query("select nextval($1::regclass) as id", [sequenceName])
const sourceId = item.sourceRow[key]
const targetId = allocated.rows[0].id
item.sourceRow[key] = targetId
if (!idMaps.has(item.table)) idMaps.set(item.table, new Map())
idMaps.get(item.table)!.set(sourceId, targetId)
}
for (const item of selected) {
const tableMetadata = metadata[item.table]
for (const foreignKey of tableMetadata.foreignKeys) {
const mapping = idMaps.get(foreignKey.referencedTable)
if (mapping?.has(item.sourceRow[foreignKey.column])) {
item.sourceRow[foreignKey.column] = mapping.get(item.sourceRow[foreignKey.column])
}
}
}
const selectedByTable = new Map<string, typeof selected>()
for (const item of selected) {
const rows = selectedByTable.get(item.table) || []
rows.push(item)
selectedByTable.set(item.table, rows)
}
const result: TenantMergeExecutionResult = {
imported: 0,
updated: 0,
retained: plan.items.length - selected.length,
remappedIds: Array.from(idMaps.values()).reduce((sum, map) => sum + Array.from(map).filter(([source, target]) => source !== target).length, 0),
tables: {},
}
for (const table of topologicalTableOrder(Array.from(selectedByTable.keys()), metadata)) {
const tableMetadata = metadata[table]
result.tables[table] ||= { imported: 0, updated: 0, retained: plan.items.filter((item) => item.table === table && !selected.includes(item)).length }
for (const item of selectedByTable.get(table) || []) {
const row = { ...item.sourceRow }
if (table === "tenants") {
delete row.locked
delete row.locked_by_export_job_id
delete row.lockedByExportJobId
}
const columns = Object.keys(row).filter((column) =>
tableMetadata.columns.includes(column) && !tableMetadata.generatedColumns.has(column)
)
const values = columns.map((column) => prepareValue(row[column], tableMetadata.jsonColumns.has(column)))
const placeholders = columns.map((_, index) => `$${index + 1}`).join(", ")
if (item.kind === "conflict" && item.targetRow) {
const primaryKey = tableMetadata.primaryKey
const updateColumns = columns.filter((column) => !primaryKey.includes(column))
if (!primaryKey.length || !updateColumns.length) continue
const whereValues = primaryKey.map((column) => item.targetRow![column])
const assignments = updateColumns.map((column) => `${quoteIdent(column)} = $${columns.indexOf(column) + 1}`).join(", ")
const where = primaryKey.map((column, index) => `${quoteIdent(column)} = $${columns.length + index + 1}`).join(" and ")
await client.query(`update ${quoteIdent(table)} set ${assignments} where ${where}`, [...values, ...whereValues])
result.updated += 1
result.tables[table].updated += 1
} else {
const inserted = await client.query(
`insert into ${quoteIdent(table)} (${columns.map(quoteIdent).join(", ")}) values (${placeholders}) on conflict do nothing`,
values
)
if (!inserted.rowCount) throw new Error(`Datensatz in ${table} konnte wegen eines neuen Konflikts nicht importiert werden`)
result.imported += 1
result.tables[table].imported += 1
}
}
}
const sourceTenant = preparedExport.tables.tenants?.find((row) => Number(row.id) === targetTenantId)
if (sourceTenant?.numberRanges) {
const targetTenantResult = await client.query(`select "numberRanges" from "tenants" where "id" = $1`, [targetTenantId])
const targetRanges = targetTenantResult.rows[0]?.numberRanges || {}
const mergedRanges = { ...targetRanges }
for (const [key, sourceRange] of Object.entries(sourceTenant.numberRanges as Record<string, any>)) {
const targetRange = targetRanges[key]
mergedRanges[key] = targetRange
? {
...sourceRange,
...targetRange,
nextNumber: Math.max(Number(sourceRange?.nextNumber || 0), Number(targetRange?.nextNumber || 0)),
}
: sourceRange
}
await client.query(`update "tenants" set "numberRanges" = $1::jsonb where "id" = $2`, [JSON.stringify(mergedRanges), targetTenantId])
}
for (const table of selectedByTable.keys()) {
const tableMetadata = metadata[table]
if (!tableMetadata.columns.includes("id")) continue
const sequenceResult = await client.query("select pg_get_serial_sequence($1, $2) as sequence_name", [`public.${table}`, "id"])
const sequenceName = sequenceResult.rows[0]?.sequence_name
if (!sequenceName) continue
await client.query(`select setval($1::regclass, greatest(coalesce((select max(id) from ${quoteIdent(table)}), 1), 1), true)`, [sequenceName])
}
await client.query("commit")
return result
} catch (err) {
await client.query("rollback")
throw err
}
}
export const executeTenantMerge = async (
exportData: TenantFullExport,
targetTenantId: number,
decisions: Record<string, TenantMergeDecision>
) => {
const client = await pool.connect()
try {
return await executeTenantMergeWithClient(client, exportData, targetTenantId, decisions)
} finally {
client.release()
}
}