diff --git a/packages/db/src/backup-lib.test.ts b/packages/db/src/backup-lib.test.ts index 13c5edc9d8..df33e46f3a 100644 --- a/packages/db/src/backup-lib.test.ts +++ b/packages/db/src/backup-lib.test.ts @@ -5,6 +5,10 @@ import { gunzipSync } from "node:zlib"; import { afterEach, describe, expect, it } from "vitest"; import postgres from "postgres"; import { createBufferedTextFileWriter, runDatabaseBackup, runDatabaseRestore } from "./backup-lib.js"; +import { + STORAGE_TRANSACTION_LOCK_DIR_NAME, + STORAGE_TRANSACTION_PARTICIPATION_NAME, +} from "./backup-transaction-lock.js"; import { ensurePostgresDatabase } from "./client.js"; import { getEmbeddedPostgresTestSupport, @@ -74,6 +78,40 @@ describe("createBufferedTextFileWriter", () => { }); }); +describe("runDatabaseBackup transaction teardown", () => { + it("does not publish an owner when database client setup throws synchronously", async () => { + const backupDir = createTempDir("paperclip-db-backup-setup-failure-"); + + await expect(runDatabaseBackup({ + connectionString: "not-a-url", + backupDir, + retention: { dailyDays: 7, weeklyWeeks: 4, monthlyMonths: 1 }, + backupEngine: "javascript", + })).rejects.toThrow(); + + expect(fs.existsSync(path.join(backupDir, STORAGE_TRANSACTION_LOCK_DIR_NAME))).toBe(false); + expect(fs.existsSync(path.join(backupDir, STORAGE_TRANSACTION_PARTICIPATION_NAME))).toBe(false); + }); + + it("releases an acquired owner when the backup connection fails", async () => { + const backupDir = createTempDir("paperclip-db-backup-connection-failure-"); + + await expect(runDatabaseBackup({ + connectionString: "postgres://paperclip:paperclip@127.0.0.1:1/paperclip", + backupDir, + retention: { dailyDays: 7, weeklyWeeks: 4, monthlyMonths: 1 }, + backupEngine: "javascript", + connectTimeoutSeconds: 1, + })).rejects.toThrow(); + + expect(fs.existsSync(path.join(backupDir, STORAGE_TRANSACTION_LOCK_DIR_NAME))).toBe(false); + expect(JSON.parse(fs.readFileSync( + path.join(backupDir, STORAGE_TRANSACTION_PARTICIPATION_NAME), + "utf8", + ))).toMatchObject({ state: "ready" }); + }); +}); + describeEmbeddedPostgres("runDatabaseBackup", () => { it( "keeps the newest backup for each retained calendar month", diff --git a/packages/db/src/backup-lib.ts b/packages/db/src/backup-lib.ts index 16dd50063f..74b936f274 100644 --- a/packages/db/src/backup-lib.ts +++ b/packages/db/src/backup-lib.ts @@ -556,7 +556,16 @@ export async function runDatabaseBackup(opts: RunDatabaseBackupOptions): Promise const excludedTableNames = normalizeTableNameSet(opts.excludeTables); const nullifiedColumnsByTable = normalizeNullifyColumnMap(opts.nullifyColumns); + let sql = postgres(opts.connectionString, { max: 1, connect_timeout: connectTimeout }); + let sqlClosed = false; + const closeSql = async () => { + if (sqlClosed) return; + sqlClosed = true; + await sql.end(); + }; + // Acquire the source-root transaction lock before any backup/prune filesystem activity. + // Client construction stays outside the hold so a synchronous setup error cannot strand it. mkdirSync(opts.backupDir, { recursive: true }); let lockHandle: StorageTransactionLockHandle; try { @@ -565,24 +574,21 @@ export async function runDatabaseBackup(opts: RunDatabaseBackupOptions): Promise operationKind: "backup", }); } catch (error) { + await closeSql(); if (error instanceof StorageTransactionLockError) { throw error; } throw error; } - let sql = postgres(opts.connectionString, { max: 1, connect_timeout: connectTimeout }); - let sqlClosed = false; - const closeSql = async () => { - if (sqlClosed) return; - sqlClosed = true; - await sql.end(); - }; const sqlFile = resolve(opts.backupDir, `${filenamePrefix}-${timestamp()}.sql`); const backupFile = `${sqlFile}.gz`; - const writer = createBufferedTextFileWriter(sqlFile); + let writer: ReturnType | null = null; + let operationError: unknown = null; try { + writer = createBufferedTextFileWriter(sqlFile); + const activeWriter = writer; if (backupEngine === "pg_dump" || (backupEngine === "auto" && canUsePgDump)) { await sql`SELECT 1`; try { @@ -592,7 +598,7 @@ export async function runDatabaseBackup(opts: RunDatabaseBackupOptions): Promise backupFile, connectTimeout, }); - await writer.abort(); + await activeWriter.abort(); const sizeBytes = statSync(backupFile).size; assertStorageTransactionLockHeld(lockHandle); const prunedCount = pruneOldBackups(opts.backupDir, retention, filenamePrefix, lockHandle); @@ -617,7 +623,7 @@ export async function runDatabaseBackup(opts: RunDatabaseBackupOptions): Promise await sql`SELECT 1`; - const emit = (line: string) => writer.emit(line); + const emit = (line: string) => activeWriter.emit(line); const emitStatement = (statement: string) => { emit(statement); emit(STATEMENT_BREAKPOINT); @@ -975,19 +981,19 @@ export async function runDatabaseBackup(opts: RunDatabaseBackupOptions): Promise const nullifiedColumns = nullifiedColumnsByTable.get(currentTableKey) ?? new Set(); if (effectiveBackupEngine !== "javascript" && nullifiedColumns.size === 0) { emit(`COPY ${qualifiedTableName} (${colNames}) FROM stdin;`); - await writer.writeRaw("\n"); + await activeWriter.writeRaw("\n"); const copySql = postgres(opts.connectionString, { max: 1, connect_timeout: connectTimeout }); try { const copyStream = await copySql .unsafe(`COPY ${qualifiedTableName} (${colNames}) TO STDOUT`) .readable(); for await (const chunk of copyStream) { - await writer.writeRaw(Buffer.isBuffer(chunk) ? chunk : Buffer.from(String(chunk))); + await activeWriter.writeRaw(Buffer.isBuffer(chunk) ? chunk : Buffer.from(String(chunk))); } } finally { await copySql.end(); } - await writer.writeRaw("\\.\n"); + await activeWriter.writeRaw("\\.\n"); emitStatementBoundary(); emit(""); continue; @@ -1004,7 +1010,7 @@ export async function runDatabaseBackup(opts: RunDatabaseBackupOptions): Promise ); emitStatement(`INSERT INTO ${qualifiedTableName} (${colNames}) VALUES (${values.join(", ")});`); } - await writer.drain(); + await activeWriter.drain(); } emit(""); } @@ -1057,7 +1063,7 @@ export async function runDatabaseBackup(opts: RunDatabaseBackupOptions): Promise emitStatement("COMMIT;"); emit(""); - await writer.close(); + await activeWriter.close(); // Compress the SQL file with gzip const sqlReadStream = createReadStream(sqlFile); @@ -1076,7 +1082,8 @@ export async function runDatabaseBackup(opts: RunDatabaseBackupOptions): Promise lockTokenPrefix: lockHandle.token.slice(0, 8), }; } catch (error) { - await writer.abort(); + operationError = error; + await writer?.abort(); if (existsSync(backupFile)) { try { unlinkSync(backupFile); } catch { /* ignore */ } } @@ -1085,11 +1092,24 @@ export async function runDatabaseBackup(opts: RunDatabaseBackupOptions): Promise } throw error; } finally { - await closeSql(); + let closeError: unknown = null; + try { + await closeSql(); + } catch (error) { + closeError = error; + } + let releaseError: unknown = null; try { releaseStorageTransactionLock(lockHandle, { participationState: "ready" }); - } catch { - // Release is best-effort on teardown; token mismatch means another owner already replaced us. + } catch (error) { + releaseError = error; + } + const teardownErrors = [operationError, closeError, releaseError].filter( + (error): error is NonNullable => error !== null, + ); + if (closeError || releaseError) { + if (teardownErrors.length === 1) throw teardownErrors[0]; + throw new AggregateError(teardownErrors, "Database backup and storage transaction teardown failed"); } } } diff --git a/server/src/index.ts b/server/src/index.ts index 723c118f55..0347102c31 100644 --- a/server/src/index.ts +++ b/server/src/index.ts @@ -32,6 +32,7 @@ import { reconcilePendingMigrationHistory, formatDatabaseBackupResult, runDatabaseBackup, + StorageTransactionLockError, authUsers, companies, companyMemberships, @@ -855,6 +856,13 @@ async function startServerWithDatabaseTeardown( ); return response; } catch (err) { + if (trigger === "scheduled" && err instanceof StorageTransactionLockError && err.code === "busy") { + logger.warn( + { backupDir: config.databaseBackupDir, trigger, reason: err.code }, + "Skipping scheduled database backup because the storage transaction lock is busy", + ); + return null; + } logger.error({ err, backupDir: config.databaseBackupDir, trigger }, `${label} database backup failed`); throw err; } finally {