Skip to content

Commit dc4f76c

Browse files
committed
fix(files): release purged file history atomically with restore
1 parent e43846c commit dc4f76c

8 files changed

Lines changed: 168 additions & 103 deletions

File tree

‎apps/docs/openapi-v2-files-audit.json‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4549,7 +4549,7 @@
45494549
"type": "integer",
45504550
"minimum": 1,
45514551
"maximum": 2147483647,
4552-
"description": "Version number, increasing by one per recorded version. Numbers are never reused, so a gap means retention removed that version.",
4552+
"description": "Version number, increasing by one per recorded version. Numbers are never reused, so a gap means an older version was removed by retention or deleted.",
45534553
"examples": [3]
45544554
},
45554555
"isCurrent": {

‎apps/sim/background/cleanup-file-versions.ts‎

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ import type { CleanupJobPayload } from '@/lib/billing/cleanup-dispatcher'
88
import type { PlanCategory } from '@/lib/billing/plan-helpers'
99
import { DEFAULT_DELETE_CHUNK_SIZE } from '@/lib/cleanup/batch-delete'
1010
import { retentionCleanupQueue } from '@/lib/cleanup/queue'
11-
import { deleteWorkspaceStorageObjects } from '@/lib/cleanup/storage-delete'
11+
import { StorageService } from '@/lib/uploads'
1212
import {
1313
FILE_VERSION_RETENTION_KEEP_LATEST,
1414
MAX_SUPERSEDED_FILE_VERSIONS,
@@ -117,10 +117,17 @@ function selectExpiredVersions(fileIds: string[], cutoff: Date, maxSuperseded: n
117117
* delete leaves its row — and the next run retries it — instead of orphaning the object.
118118
*/
119119
async function deleteVersions(rows: Array<{ id: string; key: string }>, label: string) {
120-
const failedKeys = await deleteWorkspaceStorageObjects(
121-
rows.map((row) => row.key),
122-
label
123-
)
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+
}
124131
const removable = rows.filter((row) => !failedKeys.has(row.key))
125132
let deleted = 0
126133
for (const batch of chunkArray(removable, DEFAULT_DELETE_CHUNK_SIZE)) {

‎apps/sim/background/cleanup-soft-deletes.test.ts‎

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,13 @@ vi.mock('@/lib/uploads', () => ({
8282

8383
vi.mock('@/lib/uploads/server/metadata', () => ({ deleteFileMetadata: mockDeleteFileMetadata }))
8484

85+
const { mockReleaseExpiredWorkspaceFileVersions } = vi.hoisted(() => ({
86+
mockReleaseExpiredWorkspaceFileVersions: vi.fn(),
87+
}))
88+
vi.mock('@/lib/uploads/contexts/workspace/workspace-file-versions', () => ({
89+
releaseExpiredWorkspaceFileVersions: mockReleaseExpiredWorkspaceFileVersions,
90+
}))
91+
8592
vi.mock('@/lib/workflows/utils', () => ({
8693
deduplicateWorkflowName: mockDeduplicateWorkflowName,
8794
}))
@@ -122,6 +129,39 @@ describe('cleanup soft deletes', () => {
122129
})
123130
})
124131

132+
it('releases the version history of expired workspace files before purging their objects', async () => {
133+
mockSelectRowsByIdChunks
134+
.mockResolvedValueOnce([])
135+
.mockResolvedValueOnce([])
136+
.mockResolvedValueOnce([
137+
{
138+
id: 'file-1',
139+
key: 'workspace/ws-1/file-1',
140+
workspaceId: 'ws-1',
141+
context: 'workspace',
142+
sizeBytes: 5,
143+
},
144+
{
145+
id: 'chat-file',
146+
key: 'mothership/chat-file',
147+
workspaceId: 'ws-1',
148+
context: 'mothership',
149+
sizeBytes: 3,
150+
},
151+
])
152+
153+
await runCleanupSoftDeletes(basePayload)
154+
155+
expect(mockReleaseExpiredWorkspaceFileVersions).toHaveBeenCalledWith(
156+
expect.anything(),
157+
['file-1'],
158+
expect.any(Date)
159+
)
160+
expect(mockReleaseExpiredWorkspaceFileVersions.mock.invocationCallOrder[0]).toBeLessThan(
161+
mockDeleteFiles.mock.invocationCallOrder[0]
162+
)
163+
})
164+
125165
it('keeps metadata rows whose object deletion failed', async () => {
126166
mockSelectRowsByIdChunks
127167
.mockResolvedValueOnce([])

‎apps/sim/background/cleanup-soft-deletes.ts‎

Lines changed: 6 additions & 69 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,6 @@ import {
1111
workflowMcpServer,
1212
workspaceFile,
1313
workspaceFiles,
14-
workspaceFileVersion,
1514
} from '@sim/db/schema'
1615
import { createLogger } from '@sim/logger'
1716
import { chunkArray } from '@sim/utils/helpers'
@@ -40,12 +39,12 @@ import {
4039
cleanupOwnerCondition,
4140
resolveCleanupOwnerScope,
4241
} from '@/lib/cleanup/resource-scope'
43-
import { deleteWorkspaceStorageObjects } from '@/lib/cleanup/storage-delete'
4442
import { deduplicateFolderName } from '@/lib/folders/naming'
4543
import { hardDeleteDocuments } from '@/lib/knowledge/documents/service'
4644
import type { StorageContext } from '@/lib/uploads'
4745
import { isUsingCloudStorage, StorageService } from '@/lib/uploads'
4846
import { allocateUniqueWorkspaceFileName } from '@/lib/uploads/contexts/workspace/workspace-file-manager'
47+
import { releaseExpiredWorkspaceFileVersions } from '@/lib/uploads/contexts/workspace/workspace-file-versions'
4948
import { deleteFileMetadata } from '@/lib/uploads/server/metadata'
5049
import { getWorkspaceFileSize } from '@/lib/uploads/shared/types'
5150
import { deduplicateWorkflowName } from '@/lib/workflows/utils'
@@ -224,64 +223,6 @@ async function cleanupWorkspaceFileStorage(
224223
return result
225224
}
226225

227-
/** Files whose version rows are read in one query while their objects are collected. */
228-
const VERSION_OBJECT_FILE_CHUNK_SIZE = 100
229-
230-
/**
231-
* Deletes the stored objects of every earlier version of the workspace files about to be purged,
232-
* before their version rows cascade away with the file. A file whose version objects could not all
233-
* be deleted is withheld from this purge, so its rows survive for the next run instead of leaving
234-
* those objects unreferenced. Each chunk re-reads which files are still deleted past the cutoff
235-
* right before deleting, so a file restored since selection keeps its history; the file-row delete
236-
* re-checks the same condition, so it skips that file too.
237-
*/
238-
async function cleanupWorkspaceFileVersionStorage(
239-
rows: WorkspaceFileScope['multiContextRows'],
240-
retentionDate: Date,
241-
label: string
242-
): Promise<{ rows: WorkspaceFileScope['multiContextRows']; failed: number }> {
243-
const workspaceRows = rows.filter((row) => row.context === 'workspace')
244-
if (workspaceRows.length === 0) return { rows, failed: 0 }
245-
246-
const currentKeyByFileId = new Map(workspaceRows.map((row) => [row.id, row.key]))
247-
const withheldFileIds = new Set<string>()
248-
for (const selectedFileIds of chunkArray(
249-
[...currentKeyByFileId.keys()],
250-
VERSION_OBJECT_FILE_CHUNK_SIZE
251-
)) {
252-
const stillExpired = await cleanupDb
253-
.select({ id: workspaceFiles.id })
254-
.from(workspaceFiles)
255-
.where(
256-
and(
257-
inArray(workspaceFiles.id, selectedFileIds),
258-
isNotNull(workspaceFiles.deletedAt),
259-
lt(workspaceFiles.deletedAt, retentionDate)
260-
)
261-
)
262-
const fileIds = stillExpired.map((file) => file.id)
263-
if (fileIds.length === 0) continue
264-
const versions = await cleanupDb
265-
.select({ fileId: workspaceFileVersion.fileId, key: workspaceFileVersion.key })
266-
.from(workspaceFileVersion)
267-
.where(inArray(workspaceFileVersion.fileId, fileIds))
268-
const fileIdByKey = new Map(
269-
versions
270-
.filter((version) => version.key !== currentKeyByFileId.get(version.fileId))
271-
.map((version) => [version.key, version.fileId])
272-
)
273-
for (const key of await deleteWorkspaceStorageObjects([...fileIdByKey.keys()], label)) {
274-
const fileId = fileIdByKey.get(key)
275-
if (fileId) withheldFileIds.add(fileId)
276-
}
277-
}
278-
279-
return {
280-
rows: rows.filter((row) => !withheldFileIds.has(row.id)),
281-
failed: withheldFileIds.size,
282-
}
283-
}
284-
285226
async function deleteExpiredLegacyWorkspaceFileRows(
286227
rows: WorkspaceFileScope['legacyRows'],
287228
retentionDate: Date,
@@ -959,16 +900,12 @@ export async function runCleanupSoftDeletes(
959900
chatCleanup = await prepareChatCleanup([...doomedChatIds], label)
960901
}
961902

962-
const versionCleanup = await cleanupWorkspaceFileVersionStorage(
963-
fileScope.multiContextRows,
964-
retentionDate,
965-
label
903+
await releaseExpiredWorkspaceFileVersions(
904+
cleanupDb,
905+
fileScope.multiContextRows.filter((row) => row.context === 'workspace').map((row) => row.id),
906+
retentionDate
966907
)
967-
if (budgets && versionCleanup.failed) throw new Error('File version storage cleanup failed')
968-
const fileCleanup = await cleanupWorkspaceFileStorage({
969-
...fileScope,
970-
multiContextRows: versionCleanup.rows,
971-
})
908+
const fileCleanup = await cleanupWorkspaceFileStorage(fileScope)
972909
if (budgets && fileCleanup.filesFailed) throw new Error('File storage cleanup failed')
973910

974911
let totalDeleted = 0

‎apps/sim/lib/api/contracts/v2/file-versions.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@ export const v2FileVersionSchema = z
5050
fileId: z.string().describe('File this version belongs to.'),
5151
version: versionNumberSchema
5252
.describe(
53-
'Version number, increasing by one per recorded version. Numbers are never reused, so a gap means retention removed that version.'
53+
'Version number, increasing by one per recorded version. Numbers are never reused, so a gap means an older version was removed by retention or deleted.'
5454
)
5555
.meta({ examples: [3] }),
5656
isCurrent: z.boolean().describe('Whether this version holds the current content of the file.'),

‎apps/sim/lib/cleanup/storage-delete.ts‎

Lines changed: 0 additions & 26 deletions
This file was deleted.

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

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,9 +6,11 @@ import path from 'node:path'
66
import { db, dbFor } from '@sim/db'
77
import {
88
organization,
9+
outboxEvent,
910
user,
1011
workspace,
1112
workspaceFileSecretProvenance,
13+
workspaceFiles,
1214
workspaceFileVersion,
1315
} from '@sim/db/schema'
1416
import { generateId } from '@sim/utils/id'
@@ -33,9 +35,11 @@ import {
3335
updateWorkspaceFileContent,
3436
uploadWorkspaceFile,
3537
} from '@/lib/uploads/contexts/workspace/workspace-file-manager'
38+
import { WORKSPACE_FILE_STORAGE_CLEANUP_OUTBOX_EVENT } from '@/lib/uploads/contexts/workspace/workspace-file-storage-cleanup-outbox'
3639
import {
3740
getCurrentWorkspaceFileVersion,
3841
queryWorkspaceFileVersions,
42+
releaseExpiredWorkspaceFileVersions,
3943
} from '@/lib/uploads/contexts/workspace/workspace-file-versions'
4044
import { revertWorkspaceFileVersion } from '@/lib/workspace-files/application/file-versions'
4145
import { runCleanupFileVersions } from '@/background/cleanup-file-versions'
@@ -398,6 +402,49 @@ describe('workspace file version history in PostgreSQL', () => {
398402
).rejects.toMatchObject({ code: 'conflict' })
399403
})
400404

405+
it('releases an expired file history atomically and leaves a restored file untouched', async () => {
406+
const expired = await seedFile('one')
407+
const restored = await seedFile('one')
408+
for (const fixture of [expired, restored]) {
409+
for (const content of ['two', 'three']) {
410+
await updateWorkspaceFileContent(
411+
fixture.workspaceId,
412+
fixture.fileId,
413+
fixture.aliceId,
414+
Buffer.from(content),
415+
undefined,
416+
{ version: { source: 'api', authorUserId: fixture.aliceId } }
417+
)
418+
}
419+
}
420+
const deletedAt = sql`now() - interval '40 days'`
421+
await db.update(workspaceFiles).set({ deletedAt }).where(eq(workspaceFiles.id, expired.fileId))
422+
const releasedKeys = (await versionRows(expired.fileId))
423+
.filter((row) => row.supersededAt !== null)
424+
.map((row) => row.key)
425+
426+
await releaseExpiredWorkspaceFileVersions(
427+
db,
428+
[expired.fileId, restored.fileId],
429+
new Date(Date.now() - 30 * 24 * 60 * 60 * 1000)
430+
)
431+
432+
expect((await versionRows(expired.fileId)).map((row) => row.version)).toEqual([3])
433+
expect((await versionRows(restored.fileId)).map((row) => row.version)).toEqual([1, 2, 3])
434+
const events = await db
435+
.select({ id: outboxEvent.id, payload: outboxEvent.payload })
436+
.from(outboxEvent)
437+
.where(eq(outboxEvent.eventType, WORKSPACE_FILE_STORAGE_CLEANUP_OUTBOX_EVENT))
438+
const enqueuedKeys = events.map((event) => (event.payload as { key: string }).key)
439+
expect(enqueuedKeys).toEqual(expect.arrayContaining(releasedKeys))
440+
await db.delete(outboxEvent).where(
441+
inArray(
442+
outboxEvent.id,
443+
events.map((event) => event.id)
444+
)
445+
)
446+
})
447+
401448
it('prunes superseded versions past retention but always keeps the newest ten', async () => {
402449
const fixture = await seedFile('v1')
403450
for (let index = 2; index <= 13; index++) {

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

Lines changed: 61 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,11 +3,12 @@ import {
33
type WorkspaceFileRow,
44
type WorkspaceFileVersionRow,
55
type WorkspaceFileVersionSource,
6+
workspaceFiles,
67
workspaceFileVersion,
78
} from '@sim/db/schema'
89
import { sha256Hex } from '@sim/security/hash'
910
import { generateId } from '@sim/utils/id'
10-
import { and, desc, eq, inArray, isNotNull } from 'drizzle-orm'
11+
import { and, desc, eq, inArray, isNotNull, lt, notInArray } from 'drizzle-orm'
1112
import {
1213
type CursorKey,
1314
type KeysetKey,
@@ -21,6 +22,7 @@ import {
2122
import type { DbOrTx, DbTransaction } from '@/lib/db/types'
2223
import type { WorkspaceFileRecord } from '@/lib/uploads/contexts/workspace/workspace-file-manager'
2324
import type { WorkspaceFileSecretProvenanceSnapshot } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance'
25+
import { enqueueWorkspaceFileStorageCleanups } from '@/lib/uploads/contexts/workspace/workspace-file-storage-cleanup-outbox'
2426
import { getWorkspaceFileSize } from '@/lib/uploads/shared/types'
2527

2628
/** A coalescing version stops absorbing writes this long after it was opened. */
@@ -290,6 +292,64 @@ export async function deleteWorkspaceFileVersionInTx(
290292
return deleted?.key ?? null
291293
}
292294

295+
/** Files whose history is released in one transaction. */
296+
const RELEASE_FILE_CHUNK_SIZE = 100
297+
/** Most cleanup events one outbox insert may carry. */
298+
const RELEASE_ENQUEUE_CHUNK_SIZE = 1000
299+
300+
/**
301+
* Releases the history of soft-deleted files about to be purged: per chunk, one transaction locks
302+
* the files still deleted before `deletedBefore`, deletes their superseded version rows, and
303+
* enqueues those objects on the durable storage-cleanup outbox. The row lock orders this against a
304+
* restore, which updates the same row, so a file is either restored first and keeps its whole
305+
* history or purged first and restored without it — version rows never outlive their bytes, and no
306+
* object is left unreferenced. The current version is left to the file's own purge.
307+
*/
308+
export async function releaseExpiredWorkspaceFileVersions(
309+
client: typeof db,
310+
fileIds: readonly string[],
311+
deletedBefore: Date
312+
): Promise<void> {
313+
for (let start = 0; start < fileIds.length; start += RELEASE_FILE_CHUNK_SIZE) {
314+
const chunk = fileIds.slice(start, start + RELEASE_FILE_CHUNK_SIZE)
315+
await client.transaction(async (tx) => {
316+
const expired = await tx
317+
.select({ id: workspaceFiles.id, key: workspaceFiles.key })
318+
.from(workspaceFiles)
319+
.where(
320+
and(
321+
inArray(workspaceFiles.id, chunk),
322+
isNotNull(workspaceFiles.deletedAt),
323+
lt(workspaceFiles.deletedAt, deletedBefore)
324+
)
325+
)
326+
.for('update')
327+
if (expired.length === 0) return
328+
const released = await tx
329+
.delete(workspaceFileVersion)
330+
.where(
331+
and(
332+
inArray(
333+
workspaceFileVersion.fileId,
334+
expired.map((file) => file.id)
335+
),
336+
notInArray(
337+
workspaceFileVersion.key,
338+
expired.map((file) => file.key)
339+
)
340+
)
341+
)
342+
.returning({ key: workspaceFileVersion.key })
343+
for (let offset = 0; offset < released.length; offset += RELEASE_ENQUEUE_CHUNK_SIZE) {
344+
await enqueueWorkspaceFileStorageCleanups(
345+
tx,
346+
released.slice(offset, offset + RELEASE_ENQUEUE_CHUNK_SIZE).map((row) => row.key)
347+
)
348+
}
349+
})
350+
}
351+
}
352+
293353
/** Deletes superseded versions beyond {@link MAX_SUPERSEDED_FILE_VERSIONS}, returning their keys. */
294354
async function pruneExcessWorkspaceFileVersionsInTx(
295355
tx: DbTransaction,

0 commit comments

Comments
 (0)