Skip to content

Commit 804d2da

Browse files
committed
Merge foundation commit-rejection cleanup correction
2 parents 0632de3 + bbf557f commit 804d2da

4 files changed

Lines changed: 159 additions & 9 deletions

File tree

‎apps/sim/lib/uploads/contexts/workspace/__integration__/file-versions.integration.ts‎

Lines changed: 106 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ import {
2020
import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
2121
import { sha256Hex } from '@sim/security/hash'
2222
import { createDeferred } from '@sim/testing/helpers/deferred'
23+
import { getPostgresErrorCode } from '@sim/utils/errors'
2324
import { generateId } from '@sim/utils/id'
2425
import { and, asc, eq, inArray, sql } from 'drizzle-orm'
2526
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
@@ -272,21 +273,119 @@ describe('workspace file version history in PostgreSQL', () => {
272273
}
273274
)
274275

276+
it.each(['upload', 'content'] as const)(
277+
'cleans staged %s bytes when PostgreSQL rejects COMMIT after the callback completes',
278+
async (operation) => {
279+
const fixture = await seedFile('original')
280+
const [before] = await db
281+
.select()
282+
.from(workspace)
283+
.where(eq(workspace.id, fixture.workspaceId))
284+
const triggerName = sql.identifier(`reject_commit_${generateId().replaceAll('-', '')}`)
285+
await db.execute(sql`CREATE FUNCTION ${triggerName}() RETURNS trigger LANGUAGE plpgsql AS $$
286+
BEGIN
287+
IF NEW.workspace_id = TG_ARGV[0] THEN
288+
RAISE EXCEPTION 'Deferred file constraint rejected COMMIT' USING ERRCODE = '23514';
289+
END IF;
290+
RETURN NEW;
291+
END;
292+
$$`)
293+
await db.execute(sql`CREATE CONSTRAINT TRIGGER ${triggerName}
294+
AFTER INSERT OR UPDATE ON ${workspaceFiles} DEFERRABLE INITIALLY DEFERRED
295+
FOR EACH ROW EXECUTE FUNCTION ${triggerName}(${sql.raw(`'${fixture.workspaceId}'`)})`)
296+
const transaction = db.transaction.bind(db)
297+
let callbackCompleted = false
298+
const observeCallback = vi
299+
.spyOn(db, 'transaction')
300+
.mockImplementationOnce((callback, config) =>
301+
transaction(async (tx) => {
302+
const result = await callback(tx)
303+
callbackCompleted = true
304+
return result
305+
}, config)
306+
)
307+
const upload = storageService.uploadFile
308+
let stagedKey = ''
309+
const capture = vi.spyOn(storageService, 'uploadFile').mockImplementation(async (args) => {
310+
const result = await upload(args)
311+
stagedKey = result.key
312+
return result
313+
})
314+
try {
315+
const result =
316+
operation === 'content'
317+
? updateWorkspaceFileContent(
318+
fixture.workspaceId,
319+
fixture.fileId,
320+
fixture.aliceId,
321+
Buffer.from('rejected replacement content'),
322+
undefined,
323+
{ version: { source: 'api', authorUserId: fixture.aliceId } }
324+
)
325+
: uploadWorkspaceFile(
326+
fixture.workspaceId,
327+
fixture.aliceId,
328+
Buffer.from('rejected upload'),
329+
'rejected.txt',
330+
'text/plain',
331+
{ notifyWorkspaceChange: false }
332+
)
333+
const rejection = await result.catch((error: unknown) => error)
334+
expect(getPostgresErrorCode(rejection)).toBe('23514')
335+
expect(callbackCompleted).toBe(true)
336+
} finally {
337+
capture.mockRestore()
338+
observeCallback.mockRestore()
339+
await db.execute(sql`DROP TRIGGER ${triggerName} ON ${workspaceFiles}`)
340+
await db.execute(sql`DROP FUNCTION ${triggerName}()`)
341+
}
342+
expect(stagedKey).not.toBe('')
343+
expect(
344+
await db
345+
.select({ id: workspaceFiles.id })
346+
.from(workspaceFiles)
347+
.where(eq(workspaceFiles.key, stagedKey))
348+
).toEqual([])
349+
const [after] = await db.select().from(workspace).where(eq(workspace.id, fixture.workspaceId))
350+
expect(after.storageUsedBytes).toBe(before.storageUsedBytes)
351+
expect(await versionRows(fixture.fileId)).toEqual([])
352+
const retained = await getWorkspaceFile(fixture.workspaceId, fixture.fileId)
353+
if (!retained) throw new Error('original file missing')
354+
expect((await fetchWorkspaceFileBuffer(retained, { maxBytes: 1024 })).toString()).toBe(
355+
'original'
356+
)
357+
expect(await objectExists(stagedKey)).toBe(false)
358+
expect(
359+
await db
360+
.select({ id: outboxEvent.id })
361+
.from(outboxEvent)
362+
.where(
363+
and(
364+
eq(outboxEvent.eventType, WORKSPACE_FILE_STORAGE_CLEANUP_OUTBOX_EVENT),
365+
sql`${outboxEvent.payload}->>'key' = ${stagedKey}`
366+
)
367+
)
368+
).toHaveLength(1)
369+
}
370+
)
371+
275372
it.each([
276-
{ operation: 'upload', enqueueAvailable: true },
277-
{ operation: 'content', enqueueAvailable: true },
278-
{ operation: 'upload', enqueueAvailable: false },
279-
{ operation: 'content', enqueueAvailable: false },
373+
{ operation: 'upload', enqueueAvailable: true, code: undefined },
374+
{ operation: 'content', enqueueAvailable: true, code: undefined },
375+
{ operation: 'upload', enqueueAvailable: false, code: undefined },
376+
{ operation: 'content', enqueueAvailable: false, code: undefined },
377+
{ operation: 'content', enqueueAvailable: true, code: '40003' },
378+
{ operation: 'upload', enqueueAvailable: false, code: 'CONNECTION_CLOSED' },
280379
] as const)(
281-
'retains committed $operation bytes after acknowledgement loss (cleanup database available=$enqueueAvailable)',
282-
async ({ operation, enqueueAvailable }) => {
380+
'retains committed $operation bytes after acknowledgement loss (cleanup database available=$enqueueAvailable, code=$code)',
381+
async ({ operation, enqueueAvailable, code }) => {
283382
const fixture = await seedFile('original')
284383
const transaction = db.transaction.bind(db)
285384
const lostAcknowledgement = vi
286385
.spyOn(db, 'transaction')
287386
.mockImplementationOnce(async (callback, config) => {
288387
await transaction(callback, config)
289-
throw new Error('Commit acknowledgement lost')
388+
throw Object.assign(new Error('Commit acknowledgement lost'), { code })
290389
})
291390
const enqueue = storageCleanup.enqueueWorkspaceFileStorageCleanups
292391
const enqueueFailure = vi

‎apps/sim/lib/uploads/contexts/workspace/workspace-file-manager.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import {
1919
getErrorMessage,
2020
getPostgresConstraintName,
2121
getPostgresErrorCode,
22+
isPostgresCommitRejection,
2223
} from '@sim/utils/errors'
2324
import { generateShortId } from '@sim/utils/id'
2425
import { omit } from '@sim/utils/object'
@@ -611,7 +612,7 @@ export async function discardStagedFileContent(staged: StagedFileContent): Promi
611612
})
612613
}
613614

614-
/** Keeps possibly committed bytes when the transaction fails after its callback has completed. */
615+
/** Discards known rollbacks while retaining bytes whose COMMIT outcome is uncertain. */
615616
async function finalizeStagedFileContent<T>(
616617
staged: StagedFileContent,
617618
prepare: (tx: DbTransaction) => Promise<T>
@@ -624,7 +625,7 @@ async function finalizeStagedFileContent<T>(
624625
return result
625626
})
626627
} catch (error) {
627-
if (preparedForCommit) {
628+
if (preparedForCommit && !isPostgresCommitRejection(error)) {
628629
logger.error(
629630
'File commit outcome is uncertain; retaining staged content for reconciliation',
630631
{

‎packages/utils/src/errors.test.ts‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import {
44
getPostgresCancellationReason,
55
getPostgresErrorCode,
66
getTransientDatabaseFailure,
7+
isPostgresCommitRejection,
78
} from '@sim/utils/errors'
89
import { describe, expect, it } from 'vitest'
910

@@ -236,3 +237,28 @@ describe('describeError', () => {
236237
expect(described?.causeChain?.length).toBeLessThanOrEqual(10)
237238
})
238239
})
240+
241+
describe('isPostgresCommitRejection', () => {
242+
it('recognizes a deferred constraint rejection through a transaction wrapper', () => {
243+
const rejection = Object.assign(new Error('deferred constraint failed'), { code: '23514' })
244+
expect(isPostgresCommitRejection(new Error('transaction failed', { cause: rejection }))).toBe(
245+
true
246+
)
247+
})
248+
249+
it.each(['40003', '08007', 'CONNECTION_CLOSED', 'ECONNRESET', undefined])(
250+
'does not treat uncertain outcome %s as permission to discard staged data',
251+
(code) => {
252+
const error = Object.assign(new Error('commit outcome unavailable'), { code })
253+
expect(isPostgresCommitRejection(error)).toBe(false)
254+
}
255+
)
256+
257+
it('does not replace an outer connection failure with a nested rejection code', () => {
258+
const rejection = Object.assign(new Error('earlier constraint failure'), { code: '23514' })
259+
const connection = Object.assign(new Error('connection lost', { cause: rejection }), {
260+
code: 'CONNECTION_CLOSED',
261+
})
262+
expect(isPostgresCommitRejection(connection)).toBe(false)
263+
})
264+
})

‎packages/utils/src/errors.ts‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,30 @@ export function getPostgresErrorCode(error: unknown): string | undefined {
3030
return readPgErrorField(error, 'code')
3131
}
3232

33+
const POSTGRES_COMMIT_REJECTION_CODES = new Set([
34+
'23000',
35+
'23001',
36+
'23502',
37+
'23503',
38+
'23505',
39+
'23514',
40+
'23P01',
41+
'40000',
42+
'40001',
43+
'40002',
44+
'40P01',
45+
])
46+
47+
/**
48+
* Recognizes explicit PostgreSQL constraint/rollback rejections at COMMIT. Use only on the
49+
* transaction's rejection after its callback completed, never on errors from post-commit work.
50+
* Unknown outcomes, including 40003 and connection failures, are not safe cleanup authority.
51+
*/
52+
export function isPostgresCommitRejection(error: unknown): boolean {
53+
const code = getPostgresErrorCode(error)
54+
return code !== undefined && POSTGRES_COMMIT_REJECTION_CODES.has(code)
55+
}
56+
3357
const POSTGRES_CANCELLATION_REASONS = [
3458
['57014', 'canceling statement due to statement timeout', 'statement_timeout'],
3559
['57014', 'canceling statement due to user request', 'user_cancel'],

0 commit comments

Comments
 (0)