Skip to content

Commit 9aed0f6

Browse files
committed
fix(files): clean Project bytes after proven commit rejection
1 parent 804d2da commit 9aed0f6

2 files changed

Lines changed: 186 additions & 78 deletions

File tree

‎apps/sim/lib/projects/files/__integration__/content.integration.ts‎

Lines changed: 184 additions & 76 deletions
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,7 @@ import { createDeferred } from '@sim/testing/helpers/deferred'
2828
import { emailMailerMock, emailMailerMockFns } from '@sim/testing/mocks/email-mailer.mock'
2929
import { featureFlagsMock, featureFlagsMockFns } from '@sim/testing/mocks/feature-flags.mock'
3030
import { setUploadDirServer, uploadsSetupMock } from '@sim/testing/mocks/uploads-setup.mock'
31-
import { getErrorMessage } from '@sim/utils/errors'
31+
import { getErrorMessage, getPostgresErrorCode } from '@sim/utils/errors'
3232
import { sleep } from '@sim/utils/helpers'
3333
import { generateId } from '@sim/utils/id'
3434
import { isRecordLike } from '@sim/utils/object'
@@ -2285,94 +2285,202 @@ describe('Public file shares against PostgreSQL and private storage', () => {
22852285
})
22862286

22872287
describe('Project outer transaction cleanup', () => {
2288-
for (const operation of ['create', 'update', 'revert', 'preview'] as const) {
2289-
check(`retains committed Project ${operation} bytes after acknowledgement loss`, async () => {
2288+
for (const operation of ['create', 'update'] as const) {
2289+
check(`cleans Project ${operation} bytes after a deferred COMMIT rejection`, async () => {
22902290
const f = await fixture()
2291-
const sourceBytes = Buffer.alloc(16)
2292-
sourceBytes.writeUInt32BE(16, 0)
2293-
sourceBytes.write('ftypheic', 4, 'ascii')
22942291
const source = await createProjectFile.execute({
22952292
principal: f.principal,
2296-
input:
2297-
operation === 'preview'
2298-
? {
2299-
...createInput(f.projectId),
2300-
name: 'preview.heic',
2301-
contentType: 'image/heic',
2302-
content: sourceBytes.toString('base64'),
2303-
encoding: 'base64',
2304-
}
2305-
: createInput(f.projectId, 'original'),
2293+
input: createInput(f.projectId, 'original'),
23062294
})
2307-
if (operation === 'revert')
2308-
await updateProjectFileContent.execute({
2309-
principal: f.principal,
2310-
input: {
2311-
projectId: f.projectId,
2312-
fileId: source.file.id,
2313-
content: 'second',
2314-
encoding: 'utf-8',
2315-
},
2316-
})
2317-
vi.spyOn(heic, 'transcodeHeicToJpeg').mockResolvedValue(Buffer.alloc(128, 255))
2295+
const beforeUsage = await ledger(f.organizationId)
2296+
const triggerName = sql.identifier(
2297+
`reject_project_commit_${generateId().replaceAll('-', '')}`
2298+
)
2299+
await db.execute(sql`CREATE FUNCTION ${triggerName}() RETURNS trigger LANGUAGE plpgsql AS $$
2300+
BEGIN
2301+
IF NEW.project_id = TG_ARGV[0] THEN
2302+
RAISE EXCEPTION 'Deferred file constraint rejected COMMIT' USING ERRCODE = '23514';
2303+
END IF;
2304+
RETURN NEW;
2305+
END;
2306+
$$`)
2307+
await db.execute(sql`CREATE CONSTRAINT TRIGGER ${triggerName}
2308+
AFTER INSERT OR UPDATE ON ${workspaceFiles} DEFERRABLE INITIALLY DEFERRED
2309+
FOR EACH ROW EXECUTE FUNCTION ${triggerName}(${sql.raw(`'${f.projectId}'`)})`)
2310+
const transaction = db.transaction.bind(db)
2311+
let callbackCompleted = false
2312+
const observer = vi.spyOn(db, 'transaction').mockImplementation((callback, config) =>
2313+
transaction(async (tx) => {
2314+
const result = await callback(tx)
2315+
if (isRecordLike(result) && isRecordLike(result.result) && 'file' in result.result)
2316+
callbackCompleted = true
2317+
return result
2318+
}, config)
2319+
)
23182320
const upload = storage.uploadFile
2319-
const written: string[] = []
2320-
vi.spyOn(storage, 'uploadFile').mockImplementation(async (args) => {
2321+
let stagedKey = ''
2322+
const capture = vi.spyOn(storage, 'uploadFile').mockImplementation(async (args) => {
23212323
const result = await upload(args)
2322-
written.push(result.key)
2323-
return result
2324-
})
2325-
const transaction = db.transaction.bind(db)
2326-
const primary = new Error('Commit acknowledgement lost')
2327-
let lost = false
2328-
vi.spyOn(db, 'transaction').mockImplementation(async (callback, config) => {
2329-
const result = await transaction(callback, config)
2330-
if (
2331-
!lost &&
2332-
written.length > 0 &&
2333-
isRecordLike(result) &&
2334-
isRecordLike(result.result) &&
2335-
'file' in result.result
2336-
) {
2337-
lost = true
2338-
throw primary
2339-
}
2324+
stagedKey = result.key
23402325
return result
23412326
})
2342-
const input = { projectId: f.projectId, fileId: source.file.id }
2343-
const action =
2344-
operation === 'create'
2345-
? createProjectFile.execute({
2346-
principal: f.principal,
2347-
input: { ...createInput(f.projectId, 'committed'), name: 'new.md' },
2348-
})
2349-
: operation === 'update'
2350-
? updateProjectFileContent.execute({
2327+
try {
2328+
const pending =
2329+
operation === 'create'
2330+
? createProjectFile.execute({
23512331
principal: f.principal,
2352-
input: { ...input, content: 'committed', encoding: 'utf-8' },
2332+
input: { ...createInput(f.projectId, 'rejected'), name: 'new.md' },
23532333
})
2354-
: operation === 'revert'
2355-
? revertProjectFileVersion.execute({
2356-
principal: f.principal,
2357-
input: { ...input, version: 1, expectedCurrentVersion: 2 },
2358-
})
2359-
: readProjectFileArtifact.execute({
2334+
: updateProjectFileContent.execute({
2335+
principal: f.principal,
2336+
input: {
2337+
projectId: f.projectId,
2338+
fileId: source.file.id,
2339+
content: 'rejected replacement',
2340+
encoding: 'utf-8',
2341+
},
2342+
})
2343+
const rejection = await pending.catch((error: unknown) => error)
2344+
expect(getPostgresErrorCode(rejection)).toBe('23514')
2345+
expect(callbackCompleted).toBe(true)
2346+
} finally {
2347+
observer.mockRestore()
2348+
capture.mockRestore()
2349+
await db.execute(sql`DROP TRIGGER ${triggerName} ON ${workspaceFiles}`)
2350+
await db.execute(sql`DROP FUNCTION ${triggerName}()`)
2351+
}
2352+
expect(stagedKey).not.toBe('')
2353+
expect(await ledger(f.organizationId)).toBe(beforeUsage)
2354+
expect(
2355+
await db.select().from(workspaceFiles).where(eq(workspaceFiles.key, stagedKey))
2356+
).toEqual([])
2357+
expect(
2358+
await db
2359+
.select()
2360+
.from(workspaceFileVersion)
2361+
.where(eq(workspaceFileVersion.fileId, source.file.id))
2362+
).toEqual([])
2363+
expect(
2364+
(
2365+
await readProjectFileContent.execute({
2366+
principal: f.principal,
2367+
input: { projectId: f.projectId, fileId: source.file.id },
2368+
})
2369+
).content.toString()
2370+
).toBe('original')
2371+
await expect
2372+
.soft(readFile(join(localStorageRoot, stagedKey)))
2373+
.rejects.toMatchObject({ code: 'ENOENT' })
2374+
expect(
2375+
await db
2376+
.select()
2377+
.from(outboxEvent)
2378+
.where(sql`${outboxEvent.payload}->>'key' = ${stagedKey}`)
2379+
).toHaveLength(1)
2380+
})
2381+
}
2382+
2383+
for (const { operation, code } of [
2384+
{ operation: 'create', code: undefined },
2385+
{ operation: 'update', code: undefined },
2386+
{ operation: 'revert', code: undefined },
2387+
{ operation: 'preview', code: undefined },
2388+
{ operation: 'update', code: '40003' },
2389+
{ operation: 'update', code: 'CONNECTION_CLOSED' },
2390+
] as const) {
2391+
check(
2392+
`retains committed Project ${operation} bytes after acknowledgement loss (${code ?? 'uncoded'})`,
2393+
async () => {
2394+
const f = await fixture()
2395+
const sourceBytes = Buffer.alloc(16)
2396+
sourceBytes.writeUInt32BE(16, 0)
2397+
sourceBytes.write('ftypheic', 4, 'ascii')
2398+
const source = await createProjectFile.execute({
2399+
principal: f.principal,
2400+
input:
2401+
operation === 'preview'
2402+
? {
2403+
...createInput(f.projectId),
2404+
name: 'preview.heic',
2405+
contentType: 'image/heic',
2406+
content: sourceBytes.toString('base64'),
2407+
encoding: 'base64',
2408+
}
2409+
: createInput(f.projectId, 'original'),
2410+
})
2411+
if (operation === 'revert')
2412+
await updateProjectFileContent.execute({
2413+
principal: f.principal,
2414+
input: {
2415+
projectId: f.projectId,
2416+
fileId: source.file.id,
2417+
content: 'second',
2418+
encoding: 'utf-8',
2419+
},
2420+
})
2421+
vi.spyOn(heic, 'transcodeHeicToJpeg').mockResolvedValue(Buffer.alloc(128, 255))
2422+
const upload = storage.uploadFile
2423+
const written: string[] = []
2424+
vi.spyOn(storage, 'uploadFile').mockImplementation(async (args) => {
2425+
const result = await upload(args)
2426+
written.push(result.key)
2427+
return result
2428+
})
2429+
const transaction = db.transaction.bind(db)
2430+
const primary = Object.assign(
2431+
new Error('Commit acknowledgement lost'),
2432+
code ? { code } : {}
2433+
)
2434+
let lost = false
2435+
vi.spyOn(db, 'transaction').mockImplementation(async (callback, config) => {
2436+
const result = await transaction(callback, config)
2437+
if (
2438+
!lost &&
2439+
written.length > 0 &&
2440+
isRecordLike(result) &&
2441+
isRecordLike(result.result) &&
2442+
'file' in result.result
2443+
) {
2444+
lost = true
2445+
throw primary
2446+
}
2447+
return result
2448+
})
2449+
const input = { projectId: f.projectId, fileId: source.file.id }
2450+
const action =
2451+
operation === 'create'
2452+
? createProjectFile.execute({
2453+
principal: f.principal,
2454+
input: { ...createInput(f.projectId, 'committed'), name: 'new.md' },
2455+
})
2456+
: operation === 'update'
2457+
? updateProjectFileContent.execute({
23602458
principal: f.principal,
2361-
input: { ...input, preview: true, maxBytes: 1024 },
2459+
input: { ...input, content: 'committed', encoding: 'utf-8' },
23622460
})
2363-
await expect(action).rejects.toBe(primary)
2364-
expect(lost).toBe(true)
2365-
expect(written).toHaveLength(1)
2366-
for (const key of written) {
2367-
expect((await readFile(join(localStorageRoot, key))).length).toBeGreaterThan(0)
2368-
expect(
2369-
await db
2370-
.select()
2371-
.from(outboxEvent)
2372-
.where(sql`${outboxEvent.payload}::jsonb ->> 'key' = ${key}`)
2373-
).toEqual([])
2461+
: operation === 'revert'
2462+
? revertProjectFileVersion.execute({
2463+
principal: f.principal,
2464+
input: { ...input, version: 1, expectedCurrentVersion: 2 },
2465+
})
2466+
: readProjectFileArtifact.execute({
2467+
principal: f.principal,
2468+
input: { ...input, preview: true, maxBytes: 1024 },
2469+
})
2470+
await expect(action).rejects.toBe(primary)
2471+
expect(lost).toBe(true)
2472+
expect(written).toHaveLength(1)
2473+
for (const key of written) {
2474+
expect((await readFile(join(localStorageRoot, key))).length).toBeGreaterThan(0)
2475+
expect(
2476+
await db
2477+
.select()
2478+
.from(outboxEvent)
2479+
.where(sql`${outboxEvent.payload}::jsonb ->> 'key' = ${key}`)
2480+
).toEqual([])
2481+
}
23742482
}
2375-
})
2483+
)
23762484
}
23772485

23782486
for (const deleteUnavailable of [false, true]) {

‎apps/sim/lib/projects/files/application/authorized-use-case.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import type { Principal } from '@sim/auth/principal'
22
import { db } from '@sim/db'
33
import { createLogger } from '@sim/logger'
4-
import { describeError } from '@sim/utils/errors'
4+
import { describeError, isPostgresCommitRejection } from '@sim/utils/errors'
55
import {
66
type AuthorizingUseCase,
77
recordProjectedUseCaseAuditEntries,
@@ -92,7 +92,7 @@ export function defineAuthorizedProjectFileUseCase<
9292
return { context, result }
9393
})
9494
} catch (error) {
95-
if (callbackCompleted) {
95+
if (callbackCompleted && !isPostgresCommitRejection(error)) {
9696
logger.error('Project file commit outcome is uncertain; retaining prepared resources', {
9797
operation: definition.operation.id,
9898
projectId: args.input.projectId,

0 commit comments

Comments
 (0)