@@ -20,6 +20,7 @@ import {
2020import { deleteWorkspaceFixture } from '@sim/db/testing/workspace-fixtures'
2121import { sha256Hex } from '@sim/security/hash'
2222import { createDeferred } from '@sim/testing/helpers/deferred'
23+ import { getPostgresErrorCode } from '@sim/utils/errors'
2324import { generateId } from '@sim/utils/id'
2425import { and , asc , eq , inArray , sql } from 'drizzle-orm'
2526import { 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
0 commit comments