fix(SUP-3933): harden backup lock teardown and busy schedule skip

Acquire the storage transaction lock only after client construction so
synchronous setup failures never publish an owner. Release/close errors
now surface through teardown aggregation, and scheduled backups skip
cleanly when the lock is busy.
This commit is contained in:
Andrew Levine 2026-09-11 15:06:08 -04:00
parent 8380269ec1
commit 7092d533ea
3 changed files with 85 additions and 19 deletions

View File

@ -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",

View File

@ -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<typeof createBufferedTextFileWriter> | 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<string>();
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<typeof error> => error !== null,
);
if (closeError || releaseError) {
if (teardownErrors.length === 1) throw teardownErrors[0];
throw new AggregateError(teardownErrors, "Database backup and storage transaction teardown failed");
}
}
}

View File

@ -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 {