diff --git a/backend/src/routes/admin.ts b/backend/src/routes/admin.ts index fdec24e..30534d4 100644 --- a/backend/src/routes/admin.ts +++ b/backend/src/routes/admin.ts @@ -23,13 +23,15 @@ import { importTenantFullExport, importTenantFullExportArchive, readTenantFullExportArchive, + restoreTenantMergeArchiveFiles, + restoreTenantMergeInlineFiles, } from "../utils/tenantFullExport"; import type { TenantFullExport } from "../utils/tenantFullExport"; import { buildSystemStatus } from "../modules/system-status.service"; import { matrixService } from "../modules/matrix.service"; import { s3 } from "../utils/s3"; import { secrets } from "../utils/secrets"; -import { createTenantMergeDryRun } from "../utils/tenantMergeService"; +import { createTenantMergeDryRun, executeTenantMerge, sanitizeTenantMergePlan } from "../utils/tenantMergeService"; import type { TenantMergePlan } from "../utils/tenantMergePlan"; export default async function adminRoutes(server: FastifyInstance) { @@ -401,7 +403,7 @@ export default async function adminRoutes(server: FastifyInstance) { ) => { try { const exportData = await parseMergeSource(source, contentType, filename); - const report = await createTenantMergeDryRun(exportData, targetTenantId); + const report = sanitizeTenantMergePlan(await createTenantMergeDryRun(exportData, targetTenantId)); const reportBuffer = Buffer.from(JSON.stringify(report), "utf8"); await s3.send(new PutObjectCommand({ @@ -1558,7 +1560,7 @@ export default async function adminRoutes(server: FastifyInstance) { status: importJob.status, filename: importJob.filename, statusUrl: `/api/admin/tenant-exports/${importJob.id}`, - reviewUrl: `/administration/tenants/${targetTenantId}/imports/${importJob.id}`, + reviewUrl: `/administration/tenant-imports/${importJob.id}`, }); } @@ -1590,7 +1592,7 @@ export default async function adminRoutes(server: FastifyInstance) { status: importJob.status, filename: importJob.filename, statusUrl: `/api/admin/tenant-exports/${importJob.id}`, - reviewUrl: `/administration/tenants/${targetTenantId}/imports/${importJob.id}`, + reviewUrl: `/administration/tenant-imports/${importJob.id}`, }); } @@ -1656,6 +1658,116 @@ export default async function adminRoutes(server: FastifyInstance) { } }); + // ------------------------------------------------------------- + // POST /admin/tenant-imports/:import_id/execute + // ------------------------------------------------------------- + server.post("/admin/tenant-imports/:import_id/execute", async (req, reply) => { + const currentUser = await requireAdmin(req, reply); + if (!currentUser) return; + const { import_id } = req.params as { import_id: string }; + const body = req.body as { decisions?: Record }; + const decisions = body?.decisions || {}; + const [job] = await server.db + .select() + .from(tenantExportJobs) + .where(eq(tenantExportJobs.id, import_id)) + .limit(1); + + if (!job || job.operation !== "merge") return reply.code(404).send({ error: "Merge-Import nicht gefunden" }); + if (job.status !== "review" || !job.storagePath) { + return reply.code(409).send({ error: "Merge-Import ist nicht zur Ausführung bereit", status: job.status }); + } + + let locked = false; + try { + const storedReport = JSON.parse((await readS3Buffer(mergeReportPath(job.id))).toString("utf8")) as TenantMergePlan; + const unresolved = storedReport.items.filter((item) => + item.kind === "conflict" && !["source", "target"].includes(decisions[item.id]) + ); + if (unresolved.length) { + return reply.code(400).send({ + error: `Für ${unresolved.length} Konflikte fehlt eine Entscheidung`, + unresolved: unresolved.map((item) => item.id), + }); + } + + await lockTenantForJob(job.tenantId, job.id); + locked = true; + const source = await readS3Buffer(job.storagePath); + const exportData = await parseMergeSource(source, job.contentType, job.filename); + const currentReport = sanitizeTenantMergePlan(await createTenantMergeDryRun(exportData, job.tenantId)); + const comparable = (report: TenantMergePlan) => report.items.map((item) => ({ + id: item.id, + kind: item.kind, + targetRow: item.table === "tenants" && item.targetRow + ? Object.fromEntries(Object.entries(item.targetRow).filter(([column]) => ![ + "locked", + "locked_by_export_job_id", + "lockedByExportJobId", + "updated_at", + "updatedAt", + ].includes(column))) + : item.targetRow, + })); + + if (JSON.stringify(comparable(currentReport)) !== JSON.stringify(comparable(storedReport))) { + const reportBuffer = Buffer.from(JSON.stringify(currentReport), "utf8"); + await s3.send(new PutObjectCommand({ + Bucket: secrets.S3_BUCKET, + Key: mergeReportPath(job.id), + Body: reportBuffer, + ContentType: "application/json", + ContentLength: reportBuffer.length, + })); + return reply.code(409).send({ + error: "Der Ziel-Tenant hat sich seit dem Dry-Run geändert. Der Bericht wurde aktualisiert.", + reviewRequired: true, + }); + } + + await server.db + .update(tenantExportJobs) + .set({ status: "running", error: null, updatedAt: new Date() }) + .where(eq(tenantExportJobs.id, job.id)); + + const result = await executeTenantMerge(exportData, job.tenantId, decisions); + const selectedFileRefs = new Set(); + for (const item of currentReport.items) { + const decision = decisions[item.id] || item.defaultDecision; + if (decision !== "source" || item.kind === "existing") continue; + if (item.table === "files" && item.sourceRow.id) selectedFileRefs.add(String(item.sourceRow.id)); + if (item.table === "letterheads" && item.sourceRow.id) selectedFileRefs.add(`letterhead:${item.sourceRow.id}`); + } + const files = job.contentType.includes("json") || job.filename.toLowerCase().endsWith(".json") + ? await restoreTenantMergeInlineFiles(exportData, job.tenantId, selectedFileRefs) + : await restoreTenantMergeArchiveFiles(source, exportData, job.tenantId, selectedFileRefs); + + await completeImportedTenantAccess(currentUser, { tenantId: job.tenantId }); + await server.db + .update(tenantExportJobs) + .set({ + status: "ready", + filesDone: 1, + filesTotal: 1, + completedAt: new Date(), + updatedAt: new Date(), + error: null, + }) + .where(eq(tenantExportJobs.id, job.id)); + + return { success: true, importId: job.id, tenantId: job.tenantId, result, files }; + } catch (err: any) { + server.log.error({ err, importId: job.id }, "Tenant-Merge-Ausführung fehlgeschlagen"); + await server.db + .update(tenantExportJobs) + .set({ status: "review", error: err?.message || String(err), updatedAt: new Date() }) + .where(eq(tenantExportJobs.id, job.id)); + return reply.code(500).send({ error: err?.message || "Merge-Import fehlgeschlagen" }); + } finally { + if (locked) await unlockTenantForJob(job.tenantId, job.id); + } + }); + // ------------------------------------------------------------- // PUT /admin/users/:user_id/access // ------------------------------------------------------------- diff --git a/backend/src/utils/tenantFullExport.ts b/backend/src/utils/tenantFullExport.ts index 523099b..f62fce6 100644 --- a/backend/src/utils/tenantFullExport.ts +++ b/backend/src/utils/tenantFullExport.ts @@ -712,6 +712,12 @@ const prepareCommunicationRoomsForImport = (exportData: TenantFullExport) => { } } +export const prepareTenantFullExportRowsForImport = (exportData: TenantFullExport) => { + encryptEntityBankAccountRowsForImport(exportData) + prepareCommunicationRoomsForImport(exportData) + return exportData +} + const cleanupImportedCommunicationRooms = async (client: any, exportData: TenantFullExport) => { const rows = exportData.tables.communication_rooms || [] if (!rows.length) return 0 @@ -900,8 +906,7 @@ export const importTenantFullExport = async ( } const exportData = remapTenantScopedExport(rawExportData, options.targetTenantId) - encryptEntityBankAccountRowsForImport(exportData) - prepareCommunicationRoomsForImport(exportData) + prepareTenantFullExportRowsForImport(exportData) const client = await pool.connect() const importOrder = [ "tenants", @@ -1099,3 +1104,53 @@ export const readTenantFullExportArchive = async (archiveBuffer: Buffer): Promis await reader.close() } } + +export const restoreTenantMergeArchiveFiles = async ( + archiveBuffer: Buffer, + rawExportData: TenantFullExport, + targetTenantId: number, + selectedFileRefs: ReadonlySet +) => { + const reader = new ZipReader(new BlobReader(new Blob([archiveBuffer]))) + try { + const entries = await reader.getEntries() + const entriesByName = new Map(entries.map((entry: any) => [entry.filename, entry])) + const manifest = JSON.parse(await readZipTextEntry(entriesByName, "manifest.json")) as TenantArchiveManifest + const exportData = remapTenantScopedExport(rawExportData, targetTenantId) + const sourcePrefix = `${rawExportData.tenantId}/` + const targetPrefix = `${targetTenantId}/` + const filteredFiles = (manifest.files || []) + .filter((file) => selectedFileRefs.has(String(file.id))) + .map((file) => ({ + ...file, + path: file.path?.startsWith(sourcePrefix) + ? `${targetPrefix}${file.path.slice(sourcePrefix.length)}` + : file.path, + })) + const filteredManifest: TenantArchiveManifest = { + ...manifest, + tenantId: targetTenantId, + files: filteredFiles, + } + const filteredExportData: TenantFullExport = { + ...exportData, + files: exportData.files.filter((file) => selectedFileRefs.has(String(file.id))), + } + + return await restoreArchiveFiles(entriesByName, filteredExportData, filteredManifest) + } finally { + await reader.close() + } +} + +export const restoreTenantMergeInlineFiles = async ( + rawExportData: TenantFullExport, + targetTenantId: number, + selectedFileRefs: ReadonlySet +) => { + const exportData = remapTenantScopedExport(rawExportData, targetTenantId) + return await restoreFiles({ + ...exportData, + files: exportData.files.filter((file) => selectedFileRefs.has(String(file.id))), + }) +} diff --git a/backend/src/utils/tenantMergePlan.ts b/backend/src/utils/tenantMergePlan.ts index 4a40c8a..a4391a3 100644 --- a/backend/src/utils/tenantMergePlan.ts +++ b/backend/src/utils/tenantMergePlan.ts @@ -33,6 +33,9 @@ const ignoredComparisonColumns = new Set([ "updatedAt", "updated_by", "updatedBy", + "locked", + "locked_by_export_job_id", + "lockedByExportJobId", ]) const normalizeText = (value: unknown) => String(value ?? "").trim().toLocaleLowerCase("de") diff --git a/backend/src/utils/tenantMergeService.ts b/backend/src/utils/tenantMergeService.ts index 0c99102..f686ee9 100644 --- a/backend/src/utils/tenantMergeService.ts +++ b/backend/src/utils/tenantMergeService.ts @@ -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 + generatedColumns: Set + 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 = {} 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 | 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, columns: string[]) => + columns.map((column) => row[column]).join("\0") + +const topologicalTableOrder = (tables: string[], metadata: Record) => { + 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 +} + +export const executeTenantMergeWithClient = async ( + client: any, + rawExportData: TenantFullExport, + targetTenantId: number, + decisions: Record +): Promise => { + 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[]> = {} + 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>() + 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() + 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)) { + 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 +) => { + const client = await pool.connect() + try { + return await executeTenantMergeWithClient(client, exportData, targetTenantId, decisions) + } finally { + client.release() + } +} diff --git a/backend/tests/tenantMergePlan.test.ts b/backend/tests/tenantMergePlan.test.ts index 61b35cb..a91e069 100644 --- a/backend/tests/tenantMergePlan.test.ts +++ b/backend/tests/tenantMergePlan.test.ts @@ -54,3 +54,15 @@ test("uses target as the safe default for two-way conflicts", () => { assert.equal(plan.items[0].defaultDecision, "target") assert.deepEqual(plan.items[0].differences.sort(), ["id", "label"]) }) + +test("ignores maintenance locks and audit timestamps during comparison", () => { + const plan = buildTenantMergePlan({ + tenants: [{ id: 42, name: "Tenant", locked: null, updatedAt: "before" }], + }, { + tenants: [{ id: 42, name: "Tenant", locked: "maintenance_tenant", updatedAt: "after" }], + }, { + tenants: { primaryKey: ["id"] }, + }) + + assert.equal(plan.items[0].kind, "existing") +})