Skip to content

Commit 0ccd2f1

Browse files
fix(cleanup): guard storage retries by file generation
1 parent a18eec1 commit 0ccd2f1

11 files changed

Lines changed: 196 additions & 64 deletions

File tree

‎apps/sim/background/cleanup-bounded.test.ts‎

Lines changed: 7 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -21,13 +21,11 @@ const { storage, prepareChat, executeChat, hardDelete, billing, decrement, reRoo
2121
})
2222
)
2323
const outbox = vi.hoisted(() => new Map<string, { eventType: string; payload: unknown }>())
24-
vi.mock('@/lib/core/outbox/service', async (importOriginal) => {
25-
const actual = await importOriginal<typeof import('@/lib/core/outbox/service')>()
24+
vi.mock('@/lib/core/outbox/service', () => {
2625
return {
27-
...actual,
28-
enqueueOutboxEvent: vi.fn(async (...args: Parameters<typeof actual.enqueueOutboxEvent>) => {
29-
const id = await actual.enqueueOutboxEvent(...args)
30-
outbox.set(id, { eventType: args[1], payload: args[2] })
26+
enqueueOutboxEvent: vi.fn(async (_tx: unknown, eventType: string, payload: unknown) => {
27+
const id = `test-outbox-${outbox.size}`
28+
outbox.set(id, { eventType, payload })
3129
return id
3230
}),
3331
processOutboxEventById: vi.fn(async (id: string, handlers: OutboxHandlerRegistry) => {
@@ -55,7 +53,7 @@ vi.mock('@/background/cleanup-soft-deletes', () => ({
5553
}))
5654
vi.mock('@/lib/uploads', () => ({
5755
isUsingCloudStorage: () => true,
58-
StorageService: { deleteFiles: storage },
56+
StorageService: { deleteFile: storage },
5957
}))
6058
vi.mock('@/lib/cleanup/chat-cleanup', () => ({ prepareChatCleanup: prepareChat }))
6159
vi.mock('@/lib/knowledge/documents/service', () => ({ hardDeleteDocuments: hardDelete }))
@@ -88,7 +86,7 @@ beforeEach(() => {
8886
vi.clearAllMocks()
8987
resetDbChainMock()
9088
outbox.clear()
91-
storage.mockResolvedValue({ deleted: 1, failed: [] })
89+
storage.mockResolvedValue(undefined)
9290
prepareChat.mockResolvedValue({ execute: executeChat })
9391
})
9492

@@ -154,7 +152,7 @@ describe('requested cleanup stages', () => {
154152
it('records committed log deletion if attached storage cleanup fails', async () => {
155153
queueTableRows(schemaMock.workflowExecutionLogs, [{ id: 'one', files: [{ key: 'blob' }] }])
156154
dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'one', files: [{ key: 'blob' }] }])
157-
storage.mockResolvedValue({ deleted: 0, failed: [{ key: 'blob', error: 'unavailable' }] })
155+
storage.mockRejectedValue(new Error('unavailable'))
158156
const run = control('workflowLogs', false)
159157
await expect(runBoundedLogScope(scope, run)).rejects.toThrow('storage cleanup is incomplete')
160158
expect(dbChainMockFns.delete).toHaveBeenCalledOnce()

‎apps/sim/executor/handlers/variables/variables-handler.test.ts‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -13,8 +13,7 @@ const { mockUploadFile } = vi.hoisted(() => ({
1313
mockUploadFile: vi.fn(),
1414
}))
1515

16-
vi.mock('@/lib/execution/payloads/large-value-metadata', async (importOriginal) => ({
17-
...(await importOriginal<typeof import('@/lib/execution/payloads/large-value-metadata')>()),
16+
vi.mock('@/lib/execution/payloads/large-value-metadata', () => ({
1817
registerLargeValueOwner: vi.fn().mockResolvedValue(true),
1918
addLargeValueReference: vi.fn().mockResolvedValue(undefined),
2019
}))

‎apps/sim/executor/orchestrators/loop.test.ts‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,8 +22,7 @@ const mockLogger =
2222
vi.mocked(createLogger).mock.calls.findIndex(([name]) => name === 'LoopOrchestrator')
2323
].value
2424

25-
vi.mock('@/lib/execution/payloads/large-value-metadata', async (importOriginal) => ({
26-
...(await importOriginal<typeof import('@/lib/execution/payloads/large-value-metadata')>()),
25+
vi.mock('@/lib/execution/payloads/large-value-metadata', () => ({
2726
registerLargeValueOwner: vi.fn().mockResolvedValue(true),
2827
addLargeValueReference: vi.fn().mockResolvedValue(undefined),
2928
}))

‎apps/sim/executor/variables/resolvers/block.test.ts‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,7 @@ import { navigatePathAsync } from '@/executor/variables/resolvers/reference-asyn
77
import { BlockResolver } from './block'
88
import { RESOLVED_EMPTY, type ResolutionContext } from './reference'
99

10-
vi.mock('@/lib/execution/payloads/large-value-metadata', async (importOriginal) => ({
11-
...(await importOriginal<typeof import('@/lib/execution/payloads/large-value-metadata')>()),
10+
vi.mock('@/lib/execution/payloads/large-value-metadata', () => ({
1211
registerLargeValueOwner: vi.fn().mockResolvedValue(true),
1312
addLargeValueReference: vi.fn().mockResolvedValue(undefined),
1413
}))

‎apps/sim/lib/cleanup/bounded-cleanup.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ Limits count **selected roots**, including roots restored before deletion. They
5353
- Bounded dispatch sets one attempt and a 180-second hard maximum. The worker stops starting new root batches after 120 seconds; cancellable child preparation also observes that work deadline.
5454
- Cleanup SQL uses transaction-local **500ms lock_timeout** and **5s statement_timeout**. The billable-file delete and storage decrement remain atomic. Orphan knowledge-base storage cleanup reuses the binding lock held by document creation, with a 15-second storage deadline. Large-value reference writers lock the value before registering references; cleanup locks and rechecks it before claiming a tombstone.
5555
- Trigger `cleanup` metadata reports each requested type's `selected`, `deleted`, `skipped`, `filesDeleted`, and `filesFailed`, plus stage, duration, and stop reason: `budgets_exhausted`, `scopes_exhausted`, `time_budget`, or `failed`.
56-
- A failure stops subsequent stages, preserves completed progress, and fails the run. External deletion is not transactional with Postgres. A hard process termination may leave the last metadata checkpoint behind actual effects. Do not blindly replay failed runs: inspect the failed stage and any external work already completed. Log and file storage is removed only for rows returned by the guarded delete. Their keys, and claimed large-value keys, enter `retention.storage.cleanup` outbox events in the same database transaction. The run attempts those events immediately; the existing outbox worker retries failures without selecting more roots. Inspect pending/dead-letter events when a run fails. Chat backend/storage cleanup still has its existing post-parent-delete recovery gap.
56+
- A failure stops subsequent stages, preserves completed progress, and fails the run. External deletion is not transactional with Postgres. A hard process termination may leave the last metadata checkpoint behind actual effects. Do not blindly replay failed runs: inspect the failed stage and any external work already completed. Log and file storage is removed only for rows returned by the guarded delete. Their keys, and claimed large-value keys, enter `retention.storage.cleanup` outbox events in the same database transaction. Each event captures metadata IDs, content versions, and deletion state; retries skip changed or newly registered bindings and lock matching rows through a cancellable 15-second storage deletion. The run attempts those events immediately; the existing outbox worker retries failures without selecting more roots. Inspect pending/dead-letter events when a run fails. Chat backend/storage cleanup still has its existing post-parent-delete recovery gap.
5757

5858
## Rollout
5959

‎apps/sim/lib/cleanup/bounded.postgres.integration.ts‎

Lines changed: 82 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -13,8 +13,8 @@ import {
1313
import { processOutboxEventById } from '@/lib/core/outbox/service'
1414
import { lockLargeValueKeysForReference } from '@/lib/execution/payloads/large-value-lock'
1515

16-
const { deleteStorageFiles } = vi.hoisted(() => ({ deleteStorageFiles: vi.fn() }))
17-
vi.mock('@/lib/uploads', () => ({ StorageService: { deleteFiles: deleteStorageFiles } }))
16+
const { deleteStorageFile } = vi.hoisted(() => ({ deleteStorageFile: vi.fn() }))
17+
vi.mock('@/lib/uploads', () => ({ StorageService: { deleteFile: deleteStorageFile } }))
1818

1919
const url = new URL(process.env.DATABASE_URL ?? '')
2020
if (url.hostname !== '127.0.0.1' || url.pathname !== '/bounded_cleanup_test') {
@@ -43,7 +43,7 @@ beforeAll(async () => {
4343
sql`CREATE TABLE execution_large_values (key text PRIMARY KEY, workspace_id text, deleted_at timestamp)`
4444
)
4545
await db.execute(
46-
sql`CREATE TABLE workspace_files (key text PRIMARY KEY, context text, deleted_at timestamp)`
46+
sql`CREATE TABLE workspace_files (id text PRIMARY KEY DEFAULT gen_random_uuid()::text, key text NOT NULL, context text, deleted_at timestamp, content_updated_at timestamp NOT NULL DEFAULT now())`
4747
)
4848
await db.execute(sql`CREATE TABLE workflow_execution_logs (execution_id text)`)
4949
await db.execute(sql`CREATE TABLE paused_executions (execution_id text, status text)`)
@@ -56,7 +56,7 @@ beforeAll(async () => {
5656
)
5757
})
5858
beforeEach(async () => {
59-
deleteStorageFiles.mockReset()
59+
deleteStorageFile.mockReset()
6060
await db.execute(sql`TRUNCATE outbox_event`)
6161
await db.execute(
6262
sql`TRUNCATE execution_large_value_references, execution_large_value_dependencies, execution_large_values, workflow_execution_logs, paused_executions, workspace_files`
@@ -269,7 +269,10 @@ describe('large-value reference lifecycle locks', () => {
269269
const table = kind === 'metadata' ? 'execution_large_values' : 'workspace_files'
270270
if (kind === 'metadata')
271271
await db.execute(sql`INSERT INTO execution_large_values VALUES ('live-key','one',NULL)`)
272-
else await db.execute(sql`INSERT INTO workspace_files VALUES ('live-key','execution',NULL)`)
272+
else
273+
await db.execute(
274+
sql`INSERT INTO workspace_files (key,context,deleted_at) VALUES ('live-key','execution',NULL)`
275+
)
273276
await db.transaction(async (tx) => {
274277
await lockLargeValueKeysForReference(tx, ['live-key'])
275278
await expect(
@@ -310,43 +313,104 @@ describe('durable retention storage cleanup', () => {
310313
).rejects.toThrow('abort transaction')
311314
expect(await count('roots')).toBe(6)
312315
expect((await db.execute(sql`SELECT id FROM outbox_event`)).length).toBe(0)
313-
expect(deleteStorageFiles).not.toHaveBeenCalled()
316+
expect(deleteStorageFile).not.toHaveBeenCalled()
314317
})
315318
it('retries storage from the outbox after the root is gone', async () => {
316319
const [eventId] = await cleanupQuery(async (tx) => {
317320
await tx.delete(roots).where(eq(roots.id, 'a'))
318321
return enqueueRetentionStorageCleanup(tx, ['blob'], 'execution', 1)
319322
})
320-
deleteStorageFiles.mockResolvedValueOnce({
321-
deleted: 0,
322-
failed: [{ key: 'blob', error: 'offline' }],
323-
})
323+
deleteStorageFile.mockRejectedValueOnce(new Error('offline'))
324324
await expect(
325325
processRetentionStorageCleanup(control(1), 'workflows', [eventId])
326326
).rejects.toThrow('incomplete: pending')
327327
expect(await count('roots')).toBe(5)
328-
const [pending] = await db.execute<{ status: string; payload: { keys: string[] } }>(
329-
sql`SELECT status, payload FROM outbox_event WHERE id = ${eventId}`
330-
)
328+
const [pending] = await db.execute<{
329+
status: string
330+
payload: { files: Array<{ key: string }> }
331+
}>(sql`SELECT status, payload FROM outbox_event WHERE id = ${eventId}`)
331332
expect(pending.status).toBe('pending')
332-
expect(pending.payload.keys).toEqual(['blob'])
333+
expect(pending.payload.files.map((file) => file.key)).toEqual(['blob'])
333334
await db.execute(
334335
sql`UPDATE outbox_event SET available_at = now() - interval '1 second' WHERE id = ${eventId}`
335336
)
336-
deleteStorageFiles.mockResolvedValueOnce({ deleted: 1, failed: [] })
337+
deleteStorageFile.mockResolvedValueOnce(undefined)
337338
await expect(processOutboxEventById(eventId, retentionStorageOutboxHandlers)).resolves.toBe(
338339
'completed'
339340
)
340-
expect(deleteStorageFiles).toHaveBeenCalledTimes(2)
341+
expect(deleteStorageFile).toHaveBeenCalledTimes(2)
341342
expect(await count('roots')).toBe(5)
342343
})
344+
it.each(['restore', 'replace', 'new-binding'] as const)(
345+
'does not delete a newer live binding on retry: %s',
346+
async (change) => {
347+
if (change !== 'new-binding') {
348+
await db.execute(sql`INSERT INTO workspace_files(id,key,context,deleted_at)
349+
VALUES ('original','blob','execution',now() - interval '1 day')`)
350+
}
351+
const [eventId] = await cleanupQuery((tx) =>
352+
enqueueRetentionStorageCleanup(tx, ['blob'], 'execution', 1, true)
353+
)
354+
deleteStorageFile.mockRejectedValueOnce(new Error('offline'))
355+
await expect(processOutboxEventById(eventId, retentionStorageOutboxHandlers)).resolves.toBe(
356+
'pending'
357+
)
358+
if (change === 'restore') {
359+
await db.execute(sql`UPDATE workspace_files SET deleted_at = NULL,
360+
content_updated_at = content_updated_at + interval '1 second' WHERE id = 'original'`)
361+
} else {
362+
// A replacement may coexist with an older tombstone for the same key.
363+
await db.execute(
364+
sql`INSERT INTO workspace_files(id,key,context) VALUES ('replacement','blob','execution')`
365+
)
366+
}
367+
await db.execute(
368+
sql`UPDATE outbox_event SET available_at = now() - interval '1 second' WHERE id = ${eventId}`
369+
)
370+
await expect(processOutboxEventById(eventId, retentionStorageOutboxHandlers)).resolves.toBe(
371+
'completed'
372+
)
373+
expect(deleteStorageFile).toHaveBeenCalledTimes(1)
374+
const active = await db.execute(
375+
sql`SELECT id FROM workspace_files WHERE key = 'blob' AND deleted_at IS NULL`
376+
)
377+
expect(active).toHaveLength(1)
378+
}
379+
)
380+
it('holds the captured binding lock through storage deletion and tombstones only that identity', async () => {
381+
await db.execute(
382+
sql`INSERT INTO workspace_files(id,key,context) VALUES ('original','blob','execution')`
383+
)
384+
const [eventId] = await cleanupQuery((tx) =>
385+
enqueueRetentionStorageCleanup(tx, ['blob'], 'execution', 1, true)
386+
)
387+
deleteStorageFile.mockImplementationOnce(async () => {
388+
await expect(
389+
runOutsideTransactionContext(() =>
390+
db.transaction(async (tx) => {
391+
await tx.execute(sql`SET LOCAL lock_timeout = '100ms'`)
392+
await tx.execute(
393+
sql`UPDATE workspace_files SET content_updated_at = now() WHERE id = 'original'`
394+
)
395+
})
396+
)
397+
).rejects.toThrow()
398+
})
399+
await expect(processOutboxEventById(eventId, retentionStorageOutboxHandlers)).resolves.toBe(
400+
'completed'
401+
)
402+
expect(deleteStorageFile).toHaveBeenCalledOnce()
403+
expect(
404+
await db.execute(sql`SELECT id FROM workspace_files WHERE deleted_at IS NULL`)
405+
).toHaveLength(0)
406+
})
343407
it('bounds and deduplicates the persisted key batches', async () => {
344408
await cleanupQuery((tx) =>
345409
enqueueRetentionStorageCleanup(tx, ['a', 'b', 'a', 'c'], 'workspace', 2)
346410
)
347-
const events = await db.execute<{ payload: { keys: string[] } }>(
411+
const events = await db.execute<{ payload: { files: Array<{ key: string }> } }>(
348412
sql`SELECT payload FROM outbox_event ORDER BY created_at, id`
349413
)
350-
expect(events.map((event) => event.payload.keys.length).sort()).toEqual([1, 2])
414+
expect(events.map((event) => event.payload.files.length).sort()).toEqual([1, 2])
351415
})
352416
})

‎apps/sim/lib/cleanup/bounded.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ export async function setCleanupTimeouts(tx: Pick<CleanupTransaction, 'execute'>
1818
await tx.execute(sql`SET LOCAL statement_timeout = '5s'`)
1919
}
2020

21-
/** One short DB transaction; callers do storage/backend work outside this callback. */
21+
/** One bounded DB transaction; storage under a binding lock must use a cancellable deadline. */
2222
export const cleanupQuery: CleanupQuery = (query) =>
2323
dbFor('cleanup').transaction(async (tx) => {
2424
await setCleanupTimeouts(tx)

0 commit comments

Comments
 (0)