@@ -8,7 +8,7 @@ import type { CleanupJobPayload } from '@/lib/billing/cleanup-dispatcher'
88import type { PlanCategory } from '@/lib/billing/plan-helpers'
99import { DEFAULT_DELETE_CHUNK_SIZE } from '@/lib/cleanup/batch-delete'
1010import { retentionCleanupQueue } from '@/lib/cleanup/queue'
11- import { StorageService } from '@/lib/uploads'
11+ import { enqueueWorkspaceFileStorageCleanups } from '@/lib/uploads/contexts/workspace/workspace-file-storage-cleanup-outbox '
1212import {
1313 FILE_VERSION_RETENTION_KEEP_LATEST ,
1414 MAX_SUPERSEDED_FILE_VERSIONS ,
@@ -85,7 +85,6 @@ function selectExpiredVersions(fileIds: string[], cutoff: Date, maxSuperseded: n
8585 const ranked = cleanupDb
8686 . select ( {
8787 id : workspaceFileVersion . id ,
88- key : workspaceFileVersion . key ,
8988 supersededAt : workspaceFileVersion . supersededAt ,
9089 rank : sql < number > `row_number() over (partition by ${ workspaceFileVersion . fileId } order by ${ workspaceFileVersion . version } desc)` . as (
9190 'rank'
@@ -101,7 +100,7 @@ function selectExpiredVersions(fileIds: string[], cutoff: Date, maxSuperseded: n
101100 . as ( 'ranked' )
102101
103102 return cleanupDb
104- . select ( { id : ranked . id , key : ranked . key } )
103+ . select ( { id : ranked . id } )
105104 . from ( ranked )
106105 . where (
107106 and (
@@ -113,39 +112,34 @@ function selectExpiredVersions(fileIds: string[], cutoff: Date, maxSuperseded: n
113112}
114113
115114/**
116- * Deletes the stored objects first and only then the rows whose objects are gone, so a failed
117- * delete leaves its row — and the next run retries it — instead of orphaning the object.
115+ * Deletes the expired version rows and enqueues their stored objects on the storage-cleanup outbox
116+ * in the same transaction, so a row never outlives its release and every released object is
117+ * deleted durably — the outbox retries failures and treats an already-missing object as done.
118118 */
119- async function deleteVersions ( rows : Array < { id : string ; key : string } > , label : string ) {
120- const failedKeys = new Set < string > ( )
121- for ( const batch of chunkArray ( rows , DEFAULT_DELETE_CHUNK_SIZE ) ) {
122- const deletion = await StorageService . deleteFiles (
123- batch . map ( ( row ) => row . key ) ,
124- 'workspace'
125- )
126- for ( const { key, error } of deletion . failed ) {
127- failedKeys . add ( key )
128- logger . error ( `[${ label } ] Failed to delete file version object ${ key } ` , { error } )
129- }
130- }
131- const removable = rows . filter ( ( row ) => ! failedKeys . has ( row . key ) )
119+ async function deleteVersions ( rows : Array < { id : string } > ) {
132120 let deleted = 0
133- for ( const batch of chunkArray ( removable , DEFAULT_DELETE_CHUNK_SIZE ) ) {
134- const removed = await cleanupDb
135- . delete ( workspaceFileVersion )
136- . where (
137- and (
138- inArray (
139- workspaceFileVersion . id ,
140- batch . map ( ( row ) => row . id )
141- ) ,
142- isNotNull ( workspaceFileVersion . supersededAt )
121+ for ( const batch of chunkArray ( rows , DEFAULT_DELETE_CHUNK_SIZE ) ) {
122+ deleted += await cleanupDb . transaction ( async ( tx ) => {
123+ const removed = await tx
124+ . delete ( workspaceFileVersion )
125+ . where (
126+ and (
127+ inArray (
128+ workspaceFileVersion . id ,
129+ batch . map ( ( row ) => row . id )
130+ ) ,
131+ isNotNull ( workspaceFileVersion . supersededAt )
132+ )
143133 )
134+ . returning ( { key : workspaceFileVersion . key } )
135+ await enqueueWorkspaceFileStorageCleanups (
136+ tx ,
137+ removed . map ( ( row ) => row . key )
144138 )
145- . returning ( { id : workspaceFileVersion . id } )
146- deleted += removed . length
139+ return removed . length
140+ } )
147141 }
148- return { deleted, failed : rows . length - removable . length }
142+ return deleted
149143}
150144
151145export async function runCleanupFileVersions ( payload : CleanupJobPayload ) : Promise < void > {
@@ -163,25 +157,21 @@ export async function runCleanupFileVersions(payload: CleanupJobPayload): Promis
163157 )
164158
165159 let deleted = 0
166- let failed = 0
167160 for ( const group of chunkArray ( workspaceIds , WORKSPACES_PER_QUERY ) ) {
168161 const candidates = await selectCandidateFileIds ( group , cutoff , maxSuperseded )
169162 for ( const fileIds of chunkArray ( candidates , FILES_PER_QUERY ) ) {
170163 for ( let batch = 0 ; batch < MAX_BATCHES_PER_CHUNK ; batch ++ ) {
171164 const expired = await selectExpiredVersions ( fileIds , cutoff , maxSuperseded )
172165 if ( expired . length === 0 ) break
173- const result = await deleteVersions ( expired , label )
174- deleted += result . deleted
175- failed += result . failed
176- if ( expired . length < VERSIONS_PER_BATCH || result . deleted === 0 ) break
166+ const removed = await deleteVersions ( expired )
167+ deleted += removed
168+ if ( expired . length < VERSIONS_PER_BATCH || removed === 0 ) break
177169 }
178170 }
179171 }
180172
181173 const elapsed = ( ( Date . now ( ) - startTime ) / 1000 ) . toFixed ( 2 )
182- logger . info (
183- `[${ label } ] File version cleanup: ${ deleted } deleted, ${ failed } failed in ${ elapsed } s`
184- )
174+ logger . info ( `[${ label } ] File version cleanup: ${ deleted } released in ${ elapsed } s` )
185175}
186176
187177export const cleanupFileVersionsTask = task ( {
0 commit comments