Skip to content

Commit 57cd686

Browse files
committed
fix(files): bound affected subtrees and preserve fenced revisions
1 parent 4be4772 commit 57cd686

6 files changed

Lines changed: 240 additions & 54 deletions

File tree

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

Lines changed: 107 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,8 +29,11 @@ import {
2929
seedKnowledgeAclFixture,
3030
} from '@/lib/knowledge/__integration__/seed-source-access-fixture'
3131
import {
32+
bulkArchiveWorkspaceFileItems,
3233
createWorkspaceFileFolder,
3334
fileNameExistsInWorkspaceFolder,
35+
moveWorkspaceFileItems,
36+
updateWorkspaceFileFolder,
3437
workspaceFileNameFolderCondition,
3538
} from '@/lib/uploads/contexts/workspace/workspace-file-folder-manager'
3639
import {
@@ -71,6 +74,110 @@ describe('workspace file names in PostgreSQL', () => {
7174
return ids
7275
}
7376

77+
it.each(['update', 'move', 'archive'] as const)(
78+
'allows a small subtree %s when unrelated folders exceed the bulk limit',
79+
async (operation) => {
80+
const fixture = await seedWorkspace()
81+
const rootId = generateId()
82+
const targetId = generateId()
83+
const childId = generateId()
84+
await db.execute(sql`INSERT INTO folder (id, name, user_id, workspace_id, resource_type)
85+
SELECT ${rootId} || '-' || n, 'Unrelated ' || n, ${fixture.aliceId}, ${fixture.workspaceId}, 'file'
86+
FROM generate_series(1, 5001) n`)
87+
await db.execute(sql`INSERT INTO folder (id, name, user_id, workspace_id, resource_type)
88+
VALUES (${rootId}, 'Root', ${fixture.aliceId}, ${fixture.workspaceId}, 'file'),
89+
(${targetId}, 'Target', ${fixture.aliceId}, ${fixture.workspaceId}, 'file')`)
90+
await db.execute(sql`INSERT INTO folder (id, name, user_id, workspace_id, resource_type, parent_id)
91+
VALUES (${childId}, 'Child', ${fixture.aliceId}, ${fixture.workspaceId}, 'file', ${rootId})`)
92+
if (operation === 'update') {
93+
await updateWorkspaceFileFolder({
94+
workspaceId: fixture.workspaceId,
95+
folderId: rootId,
96+
parentId: targetId,
97+
})
98+
} else if (operation === 'move') {
99+
await moveWorkspaceFileItems({
100+
workspaceId: fixture.workspaceId,
101+
folderIds: [rootId],
102+
targetFolderId: targetId,
103+
})
104+
} else {
105+
await bulkArchiveWorkspaceFileItems({
106+
workspaceId: fixture.workspaceId,
107+
folderIds: [rootId],
108+
})
109+
}
110+
const rows = await db.execute<{ id: string; parentId: string | null; archived: boolean }>(sql`
111+
SELECT id, parent_id AS "parentId", deleted_at IS NOT NULL AS archived
112+
FROM folder WHERE id IN (${rootId}, ${childId}) ORDER BY name`)
113+
expect([...rows]).toEqual([
114+
{ id: childId, parentId: rootId, archived: operation === 'archive' },
115+
{
116+
id: rootId,
117+
parentId: operation === 'archive' ? null : targetId,
118+
archived: operation === 'archive',
119+
},
120+
])
121+
const [unrelated] = await db.execute<{
122+
count: number
123+
}>(sql`SELECT count(*)::int AS count FROM folder
124+
WHERE workspace_id = ${fixture.workspaceId} AND parent_id IS NULL AND deleted_at IS NULL
125+
AND id NOT IN (${rootId}, ${childId}, ${targetId})`)
126+
expect(unrelated.count).toBe(5001)
127+
}
128+
)
129+
130+
it('rejects a move whose actual subtree exceeds the bulk limit without changing its parent', async () => {
131+
const fixture = await seedWorkspace()
132+
const rootId = generateId()
133+
const targetId = generateId()
134+
await db.execute(sql`INSERT INTO folder (id, name, user_id, workspace_id, resource_type)
135+
VALUES (${rootId}, 'Root', ${fixture.aliceId}, ${fixture.workspaceId}, 'file'),
136+
(${targetId}, 'Target', ${fixture.aliceId}, ${fixture.workspaceId}, 'file')`)
137+
await db.execute(sql`INSERT INTO folder (id, name, user_id, workspace_id, resource_type, parent_id)
138+
SELECT ${rootId} || '-' || n, 'Child ' || n, ${fixture.aliceId}, ${fixture.workspaceId}, 'file', ${rootId}
139+
FROM generate_series(1, 5000) n`)
140+
await expect(
141+
updateWorkspaceFileFolder({
142+
workspaceId: fixture.workspaceId,
143+
folderId: rootId,
144+
parentId: targetId,
145+
})
146+
).rejects.toMatchObject({
147+
code: 'validation',
148+
message: 'File operation affects more than 5000 items',
149+
})
150+
expect([...(await db.execute(sql`SELECT parent_id FROM folder WHERE id = ${rootId}`))]).toEqual(
151+
[{ parent_id: null }]
152+
)
153+
})
154+
155+
it('rejects descendant destinations when unrelated folders exceed the bulk limit', async () => {
156+
const fixture = await seedWorkspace()
157+
const rootId = generateId()
158+
const childId = generateId()
159+
await db.execute(sql`INSERT INTO folder (id, name, user_id, workspace_id, resource_type)
160+
SELECT ${rootId} || '-' || n, 'Unrelated ' || n, ${fixture.aliceId}, ${fixture.workspaceId}, 'file'
161+
FROM generate_series(1, 5001) n`)
162+
await db.execute(sql`INSERT INTO folder (id, name, user_id, workspace_id, resource_type)
163+
VALUES (${rootId}, 'Root', ${fixture.aliceId}, ${fixture.workspaceId}, 'file')`)
164+
await db.execute(sql`INSERT INTO folder (id, name, user_id, workspace_id, resource_type, parent_id)
165+
VALUES (${childId}, 'Child', ${fixture.aliceId}, ${fixture.workspaceId}, 'file', ${rootId})`)
166+
await expect(
167+
updateWorkspaceFileFolder({
168+
workspaceId: fixture.workspaceId,
169+
folderId: rootId,
170+
parentId: childId,
171+
})
172+
).rejects.toMatchObject({
173+
code: 'validation',
174+
message: 'Cannot move a folder into one of its descendants',
175+
})
176+
expect([...(await db.execute(sql`SELECT parent_id FROM folder WHERE id = ${rootId}`))]).toEqual(
177+
[{ parent_id: null }]
178+
)
179+
})
180+
74181
function upload(workspaceId: string, userId: string, name: string, folderId?: string | null) {
75182
return uploadWorkspaceFile(workspaceId, userId, Buffer.from(name), name, 'text/plain', {
76183
folderId,

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

Lines changed: 29 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -175,6 +175,32 @@ function assertBulkAffectedItemsWithinLimit(count: number): void {
175175
}
176176
}
177177

178+
/** Loads only the affected active subtrees, with one extra row to detect the bulk limit. */
179+
async function selectAffectedFileFolders(
180+
owner: EditableFileOwner,
181+
folderIds: string[],
182+
tx: DbTransaction
183+
): Promise<Array<{ id: string; parentId: string | null }>> {
184+
if (folderIds.length === 0) return []
185+
const ownerCondition = fileFolderOwnerCondition(owner)
186+
const rows = await tx.execute<{ id: string; parentId: string | null }>(sql`
187+
WITH RECURSIVE affected AS (
188+
SELECT ${folderTable.id}, ${folderTable.parentId}
189+
FROM ${folderTable}
190+
WHERE ${ownerCondition} AND ${inArray(folderTable.id, folderIds)}
191+
AND ${folderTable.deletedAt} IS NULL
192+
UNION
193+
SELECT ${folderTable.id}, ${folderTable.parentId}
194+
FROM ${folderTable} JOIN affected ON ${folderTable.parentId} = affected.id
195+
WHERE ${ownerCondition} AND ${folderTable.deletedAt} IS NULL
196+
)
197+
SELECT id, parent_id AS "parentId" FROM affected
198+
LIMIT ${MAX_WORKSPACE_FILE_BULK_AFFECTED_ITEMS + 1}
199+
`)
200+
assertBulkAffectedItemsWithinLimit(rows.length)
201+
return [...rows]
202+
}
203+
178204
/**
179205
* Verifies every requested active file/folder belongs to this workspace before a bulk mutation.
180206
* This prevents the bulk archive primitive's workspace predicate from silently turning an
@@ -906,12 +932,7 @@ async function updateFileFolder<O extends EditableFileOwner>(
906932
throw new OrchestrationError('validation', 'Folder cannot be its own parent')
907933
await assertFileFolderTarget(params.owner, finalParentId, tx)
908934
if (params.parentId !== undefined) {
909-
const activeFolders = await tx
910-
.select({ id: folderTable.id, parentId: folderTable.parentId })
911-
.from(folderTable)
912-
.where(and(ownerCondition, isNull(folderTable.deletedAt)))
913-
.limit(MAX_WORKSPACE_FILE_BULK_AFFECTED_ITEMS + 1)
914-
assertBulkAffectedItemsWithinLimit(activeFolders.length)
935+
const activeFolders = await selectAffectedFileFolders(params.owner, [params.folderId], tx)
915936
if (
916937
finalParentId &&
917938
collectDescendantFolderIds(activeFolders, params.folderId).includes(finalParentId)
@@ -1060,15 +1081,7 @@ async function moveFileItems(
10601081
}
10611082

10621083
if (folderIds.length > 0) {
1063-
const activeFolders = await tx
1064-
.select({ id: folderTable.id, parentId: folderTable.parentId })
1065-
.from(folderTable)
1066-
.where(
1067-
and(fileFolderOwnerCondition(params.owner), isFileFolder, isNull(folderTable.deletedAt))
1068-
)
1069-
.limit(MAX_WORKSPACE_FILE_BULK_AFFECTED_ITEMS + 1)
1070-
1071-
assertBulkAffectedItemsWithinLimit(activeFolders.length)
1084+
const activeFolders = await selectAffectedFileFolders(params.owner, folderIds, tx)
10721085

10731086
const affectedFolderIds = new Set<string>()
10741087

@@ -1423,17 +1436,7 @@ async function archiveFileItems(
14231436
await acquireFileFolderMutationLock(tx, params.owner)
14241437
await assertFileItemsBelongToOwner(params, tx)
14251438

1426-
const activeFolders =
1427-
explicitFolderIds.length > 0
1428-
? await tx
1429-
.select({ id: folderTable.id, parentId: folderTable.parentId })
1430-
.from(folderTable)
1431-
.where(
1432-
and(fileFolderOwnerCondition(params.owner), isFileFolder, isNull(folderTable.deletedAt))
1433-
)
1434-
.limit(MAX_WORKSPACE_FILE_BULK_AFFECTED_ITEMS + 1)
1435-
: []
1436-
assertBulkAffectedItemsWithinLimit(activeFolders.length)
1439+
const activeFolders = await selectAffectedFileFolders(params.owner, explicitFolderIds, tx)
14371440
const descendantFolderIds = explicitFolderIds.flatMap((folderId) =>
14381441
collectDescendantFolderIds(activeFolders, folderId)
14391442
)

‎packages/db/file-folder-version-ownership.integration.ts‎

Lines changed: 22 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -24,11 +24,11 @@ describe('file folder and version ownership in PostgreSQL', () => {
2424
}
2525
}
2626

27-
async function waitForDatabaseLock(pid: number) {
27+
async function waitForDatabaseLock(pid: number, writerPid: number) {
2828
for (let attempt = 0; attempt < 100; attempt += 1) {
2929
const [state] = await sql<{ waiting: boolean }[]>`
3030
SELECT EXISTS (
31-
SELECT 1 FROM pg_stat_activity WHERE pid = ${pid} AND wait_event_type = 'Lock'
31+
SELECT 1 FROM unnest(pg_blocking_pids(${pid})) blocker WHERE blocker = ${writerPid}
3232
) AS waiting
3333
`
3434
if (state.waiting) return
@@ -338,12 +338,13 @@ describe('file folder and version ownership in PostgreSQL', () => {
338338
await sql`INSERT INTO folder (id, name, user_id, workspace_id, resource_type, project_id)
339339
VALUES ('folder-a', 'A', 'user-a', ${workspaceId}, 'file', ${projectId}),
340340
('folder-b', 'B', 'user-a', ${workspaceId}, 'file', ${projectId})`
341-
const moved = createDeferred<void>()
341+
const moved = createDeferred<number>()
342342
const release = createDeferred<void>()
343343
const firstMove = sql.begin(async (tx) => {
344344
await tx`SELECT pg_advisory_xact_lock(hashtextextended(${lockKey}, 0))`
345345
await tx`UPDATE folder SET parent_id = 'folder-b' WHERE id = 'folder-a'`
346-
moved.resolve()
346+
const [writer] = await tx<{ pid: number }[]>`SELECT pg_backend_pid() AS pid`
347+
moved.resolve(writer.pid)
347348
await release.promise
348349
})
349350
const waiting = createDeferred<number>()
@@ -360,12 +361,12 @@ describe('file folder and version ownership in PostgreSQL', () => {
360361
rejection = expect(secondMove).rejects.toMatchObject({
361362
code: isolation === 'repeatable read' ? '40001' : '23514',
362363
})
363-
await waitForDatabaseLock(await waiting.promise)
364+
await waitForDatabaseLock(await waiting.promise, await moved.promise)
364365
} finally {
365366
release.resolve()
366367
await firstMove
368+
await rejection
367369
}
368-
await rejection
369370
expect(await sql`SELECT id, parent_id FROM folder ORDER BY id`).toEqual([
370371
{ id: 'folder-a', parent_id: 'folder-b' },
371372
{ id: 'folder-b', parent_id: null },
@@ -378,12 +379,13 @@ describe('file folder and version ownership in PostgreSQL', () => {
378379
async (isolation) => {
379380
await sql`INSERT INTO workspace_files (id, context, user_id, workspace_id)
380381
VALUES ('file', 'workspace', 'user-a', 'workspace-a')`
381-
const inserted = createDeferred<void>()
382+
const inserted = createDeferred<number>()
382383
const release = createDeferred<void>()
383384
const creation = sql.begin(async (tx) => {
384385
await tx`INSERT INTO workspace_file_version (id, file_id, workspace_id)
385386
VALUES ('version', 'file', 'workspace-a')`
386-
inserted.resolve()
387+
const [writer] = await tx<{ pid: number }[]>`SELECT pg_backend_pid() AS pid`
388+
inserted.resolve(writer.pid)
387389
await release.promise
388390
})
389391
const waiting = createDeferred<number>()
@@ -399,12 +401,12 @@ describe('file folder and version ownership in PostgreSQL', () => {
399401
rejection = expect(transfer).rejects.toMatchObject({
400402
code: isolation === 'repeatable read' ? '40001' : '23514',
401403
})
402-
await waitForDatabaseLock(await waiting.promise)
404+
await waitForDatabaseLock(await waiting.promise, await inserted.promise)
403405
} finally {
404406
release.resolve()
405407
await creation
408+
await rejection
406409
}
407-
await rejection
408410
const [file] = await sql`SELECT workspace_id FROM workspace_files WHERE id = 'file'`
409411
expect(file.workspace_id).toBe('workspace-a')
410412
}
@@ -415,11 +417,12 @@ describe('file folder and version ownership in PostgreSQL', () => {
415417
async (isolation) => {
416418
await sql`INSERT INTO workspace_files (id, context, user_id, workspace_id)
417419
VALUES ('file', 'workspace', 'user-a', 'workspace-a')`
418-
const moved = createDeferred<void>()
420+
const moved = createDeferred<number>()
419421
const release = createDeferred<void>()
420422
const transfer = sql.begin(async (tx) => {
421423
await tx`UPDATE workspace_files SET workspace_id = 'workspace-b' WHERE id = 'file'`
422-
moved.resolve()
424+
const [writer] = await tx<{ pid: number }[]>`SELECT pg_backend_pid() AS pid`
425+
moved.resolve(writer.pid)
423426
await release.promise
424427
})
425428
const waiting = createDeferred<number>()
@@ -436,25 +439,26 @@ describe('file folder and version ownership in PostgreSQL', () => {
436439
rejection = expect(creation).rejects.toMatchObject({
437440
code: isolation === 'repeatable read' ? '40001' : '23514',
438441
})
439-
await waitForDatabaseLock(await waiting.promise)
442+
await waitForDatabaseLock(await waiting.promise, await moved.promise)
440443
} finally {
441444
release.resolve()
442445
await transfer
446+
await rejection
443447
}
444-
await rejection
445448
expect(await sql`SELECT id FROM workspace_file_version`).toHaveLength(0)
446449
}
447450
)
448451

449452
it.each(['read committed', 'repeatable read'])(
450453
'protects a Project from concurrent folder creation and deletion at %s isolation',
451454
async (isolation) => {
452-
const inserted = createDeferred<void>()
455+
const inserted = createDeferred<number>()
453456
const release = createDeferred<void>()
454457
const creation = sql.begin(async (tx) => {
455458
await tx`INSERT INTO folder (id, name, user_id, resource_type, project_id)
456459
VALUES ('folder', 'Docs', 'user-a', 'file', 'project-a')`
457-
inserted.resolve()
460+
const [writer] = await tx<{ pid: number }[]>`SELECT pg_backend_pid() AS pid`
461+
inserted.resolve(writer.pid)
458462
await release.promise
459463
})
460464
const waiting = createDeferred<number>()
@@ -470,12 +474,12 @@ describe('file folder and version ownership in PostgreSQL', () => {
470474
rejection = expect(deletion).rejects.toMatchObject({
471475
code: '23503',
472476
})
473-
await waitForDatabaseLock(await waiting.promise)
477+
await waitForDatabaseLock(await waiting.promise, await inserted.promise)
474478
} finally {
475479
release.resolve()
476480
await creation
481+
await rejection
477482
}
478-
await rejection
479483
expect(await sql`SELECT id FROM folder WHERE id = 'folder'`).toHaveLength(1)
480484
}
481485
)

‎packages/db/migrations/0403_file_folder_version_ownership.sql‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -116,6 +116,16 @@ DROP TRIGGER IF EXISTS workspace_files_validate_relations ON workspace_files;-->
116116
CREATE CONSTRAINT TRIGGER workspace_files_validate_relations
117117
AFTER INSERT OR UPDATE ON workspace_files DEFERRABLE INITIALLY IMMEDIATE
118118
FOR EACH ROW EXECUTE FUNCTION workspace_files_validate_relations();--> statement-breakpoint
119+
-- Relationship fences advance the row version without changing durable content or its provenance.
120+
-- Real inserts and metadata/content changes must still normalize before downstream triggers run.
121+
CREATE OR REPLACE FUNCTION workspace_file_content_version_millisecond()
122+
RETURNS trigger LANGUAGE plpgsql AS $$
123+
BEGIN
124+
IF TG_OP = 'UPDATE' AND NEW IS NOT DISTINCT FROM OLD THEN RETURN NEW; END IF;
125+
NEW.content_updated_at := date_trunc('milliseconds', NEW.content_updated_at);
126+
RETURN NEW;
127+
END;
128+
$$;--> statement-breakpoint
119129
CREATE OR REPLACE FUNCTION workspace_file_version_owner_match()
120130
RETURNS trigger LANGUAGE plpgsql AS $$
121131
DECLARE

0 commit comments

Comments
 (0)