Skip to content

Commit 08398e7

Browse files
committed
fix(lifecycle): preserve empty upload attribution through handoff
1 parent 34a82c3 commit 08398e7

3 files changed

Lines changed: 206 additions & 17 deletions

File tree

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

Lines changed: 160 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import {
1313
workspaceFiles,
1414
workspaceFileVersion,
1515
} from '@sim/db/schema'
16+
import { sha256Hex } from '@sim/security/hash'
1617
import { createDeferred } from '@sim/testing/helpers/deferred'
1718
import { generateId } from '@sim/utils/id'
1819
import { and, asc, eq, inArray, sql } from 'drizzle-orm'
@@ -240,12 +241,149 @@ describe('workspace file version history in PostgreSQL', () => {
240241
])
241242
})
242243

243-
it('keeps an empty shell unversioned until its first content write after handoff', async () => {
244+
it.each([false, true])(
245+
'preserves empty history across handoff and first-write version one (unknown=%s)',
246+
async (unknown) => {
247+
const fixture = await seedFile('')
248+
if (unknown)
249+
await db
250+
.update(workspaceFiles)
251+
.set({ contentUpdatedAt: new Date(Date.now() + 5000) })
252+
.where(eq(workspaceFiles.id, fixture.fileId))
253+
const before = await getWorkspaceFile(fixture.workspaceId, fixture.fileId)
254+
if (!before) throw new Error('file missing')
255+
const [accounting] = await db
256+
.select({ bytes: workspace.storageUsedBytes })
257+
.from(workspace)
258+
.where(eq(workspace.id, fixture.workspaceId))
259+
const historyBefore = await queryWorkspaceFileVersions(before, {
260+
sortOrder: 'asc',
261+
limit: 10,
262+
})
263+
expect(historyBefore.versions).toMatchObject([
264+
{
265+
version: 1,
266+
source: unknown ? 'unknown' : 'upload',
267+
authorUserIds: unknown ? [] : [fixture.aliceId],
268+
size: 0,
269+
},
270+
])
271+
for (const successor of [fixture.bobId, fixture.aliceId, fixture.bobId]) {
272+
await db.transaction((tx) =>
273+
handoffFileCreatorsInTx(tx, eq(workspaceFiles.id, fixture.fileId), successor)
274+
)
275+
}
276+
const after = await getWorkspaceFile(fixture.workspaceId, fixture.fileId)
277+
if (!after) throw new Error('file missing')
278+
expect(after.uploadedBy).toBe(fixture.bobId)
279+
const historyAfter = await queryWorkspaceFileVersions(after, { sortOrder: 'asc', limit: 10 })
280+
expect(historyAfter.versions).toEqual(historyBefore.versions)
281+
expect(
282+
await db
283+
.select({ bytes: workspace.storageUsedBytes })
284+
.from(workspace)
285+
.where(eq(workspace.id, fixture.workspaceId))
286+
).toEqual([accounting])
287+
expect(await objectExists(fixture.firstKey)).toBe(true)
288+
await updateWorkspaceFileContent(
289+
fixture.workspaceId,
290+
fixture.fileId,
291+
fixture.bobId,
292+
Buffer.from('first'),
293+
undefined,
294+
{
295+
version: { source: 'api', authorUserId: fixture.bobId },
296+
secretProvenancePolicy: { mode: 'replace', provenance: { status: 'unknown' } },
297+
}
298+
)
299+
expect(await versionRows(fixture.fileId)).toMatchObject([
300+
{
301+
version: 1,
302+
source: 'api',
303+
authorUserIds: [fixture.bobId],
304+
contentHash: sha256Hex('first'),
305+
restoredFromVersion: null,
306+
secretProvenanceStatus: 'unknown',
307+
secretProvenanceEntries: [],
308+
},
309+
])
310+
const [written] = await versionRows(fixture.fileId)
311+
expect(written.createdAt).toEqual(written.updatedAt)
312+
expect(written.createdAt).not.toEqual(historyBefore.versions[0].createdAt)
313+
expect(await objectExists(fixture.firstKey)).toBe(false)
314+
expect(await readVersionBytes(fixture.workspaceId, fixture.fileId, written.key)).toBe('first')
315+
}
316+
)
317+
318+
it.each(['api', 'revert'] as const)(
319+
'keeps an explicitly recorded empty %s version through handoff and a later write',
320+
async (source) => {
321+
const fixture = await seedFile('')
322+
await updateWorkspaceFileContent(
323+
fixture.workspaceId,
324+
fixture.fileId,
325+
fixture.aliceId,
326+
Buffer.alloc(0),
327+
undefined,
328+
{
329+
version: {
330+
source: 'api',
331+
authorUserId: fixture.aliceId,
332+
},
333+
}
334+
)
335+
if (source === 'revert') {
336+
await updateWorkspaceFileContent(
337+
fixture.workspaceId,
338+
fixture.fileId,
339+
fixture.aliceId,
340+
Buffer.from('intermediate'),
341+
undefined,
342+
{ version: { source: 'api', authorUserId: fixture.aliceId } }
343+
)
344+
await revertWorkspaceFileVersion.execute({
345+
principal: { kind: 'session', userId: fixture.aliceId, sessionId: generateId() },
346+
input: { fileId: fixture.fileId, assertedWorkspaceId: fixture.workspaceId, version: 1 },
347+
})
348+
}
349+
const previousRows = await versionRows(fixture.fileId)
350+
const recorded = previousRows[previousRows.length - 1]
351+
await db.transaction((tx) =>
352+
handoffFileCreatorsInTx(tx, eq(workspaceFiles.id, fixture.fileId), fixture.bobId)
353+
)
354+
await updateWorkspaceFileContent(
355+
fixture.workspaceId,
356+
fixture.fileId,
357+
fixture.bobId,
358+
Buffer.from('later'),
359+
undefined,
360+
{ version: { source: 'api', authorUserId: fixture.bobId } }
361+
)
362+
const rows = await versionRows(fixture.fileId)
363+
expect(rows.map((row) => [row.version, row.source, row.authorUserIds])).toEqual([
364+
...previousRows.map((row) => [row.version, row.source, row.authorUserIds]),
365+
[recorded.version + 1, 'api', [fixture.bobId]],
366+
])
367+
expect(rows[recorded.version - 1]).toMatchObject({
368+
key: recorded.key,
369+
contentHash: recorded.contentHash,
370+
restoredFromVersion: source === 'revert' ? 1 : null,
371+
})
372+
expect(await objectExists(recorded.key)).toBe(true)
373+
expect(await readVersionBytes(fixture.workspaceId, fixture.fileId, recorded.key)).toBe('')
374+
}
375+
)
376+
377+
it('keeps an empty storage object still referenced by a retained file after the first write', async () => {
244378
const fixture = await seedFile('')
379+
const [original] = await db
380+
.select()
381+
.from(workspaceFiles)
382+
.where(eq(workspaceFiles.id, fixture.fileId))
383+
await db.insert(workspaceFiles).values({ ...original, id: generateId(), deletedAt: new Date() })
245384
await db.transaction((tx) =>
246385
handoffFileCreatorsInTx(tx, eq(workspaceFiles.id, fixture.fileId), fixture.bobId)
247386
)
248-
expect(await versionRows(fixture.fileId)).toEqual([])
249387
await updateWorkspaceFileContent(
250388
fixture.workspaceId,
251389
fixture.fileId,
@@ -254,15 +392,19 @@ describe('workspace file version history in PostgreSQL', () => {
254392
undefined,
255393
{ version: { source: 'api', authorUserId: fixture.bobId } }
256394
)
257-
expect(await versionRows(fixture.fileId)).toMatchObject([
258-
{ version: 1, source: 'api', authorUserIds: [fixture.bobId] },
259-
])
395+
expect(await objectExists(fixture.firstKey)).toBe(true)
396+
expect((await versionRows(fixture.fileId)).map((row) => row.version)).toEqual([1])
260397
})
261398

262-
it.each(['handoff', 'write'] as const)(
263-
'serializes creator handoff with a content write when %s waits first',
264-
async (first) => {
265-
const fixture = await seedFile('original')
399+
it.each([
400+
{ first: 'handoff', content: 'original' },
401+
{ first: 'write', content: 'original' },
402+
{ first: 'handoff', content: '' },
403+
{ first: 'write', content: '' },
404+
] as const)(
405+
'serializes creator handoff with a content write ($first first, original=$content)',
406+
async ({ first, content }) => {
407+
const fixture = await seedFile(content)
266408
const ready = createDeferred<number>()
267409
const release = createDeferred<void>()
268410
const blocker = db.transaction(async (tx) => {
@@ -319,13 +461,17 @@ describe('workspace file version history in PostgreSQL', () => {
319461
await Promise.all(pending)
320462
}
321463
const rows = await versionRows(fixture.fileId)
322-
expect(rows.map((row) => [row.version, row.authorUserIds])).toEqual([
323-
[1, [fixture.aliceId]],
324-
[2, [fixture.aliceId]],
325-
])
464+
expect(rows.map((row) => [row.version, row.authorUserIds])).toEqual(
465+
content
466+
? [
467+
[1, [fixture.aliceId]],
468+
[2, [fixture.aliceId]],
469+
]
470+
: [[1, [fixture.aliceId]]]
471+
)
326472
expect(rows.filter((row) => row.supersededAt === null)).toHaveLength(1)
327473
expect(await readVersionBytes(fixture.workspaceId, fixture.fileId, rows[0].key)).toBe(
328-
'original'
474+
content || 'second'
329475
)
330476
expect((await getWorkspaceFile(fixture.workspaceId, fixture.fileId))?.uploadedBy).toBe(
331477
fixture.bobId

‎apps/sim/lib/uploads/contexts/workspace/creator-handoff.ts‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,6 @@ import {
77
loadWorkspaceFileVersionHead,
88
materializeWorkspaceFileVersionInTx,
99
} from '@/lib/uploads/contexts/workspace/workspace-file-versions'
10-
import { getWorkspaceFileSize } from '@/lib/uploads/shared/types'
1110

1211
/** Freezes implicit history under the writer lock before replacing a shared file's live creator. */
1312
export async function handoffFileCreatorsInTx(
@@ -25,7 +24,7 @@ export async function handoffFileCreatorsInTx(
2524
for (const file of files) {
2625
if (file.context !== 'workspace' || !file.workspaceId) continue
2726
const head = await loadWorkspaceFileVersionHead(file.id, tx)
28-
if (isVersionHeadCurrent(head, file) || (!head && getWorkspaceFileSize(file) === 0)) continue
27+
if (isVersionHeadCurrent(head, file)) continue
2928
const provenance = await snapshotWorkspaceFileSecretProvenanceInTx(
3029
tx,
3130
file.id,

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

Lines changed: 45 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -219,13 +219,34 @@ export async function materializeWorkspaceFileVersionInTx(
219219
return materialized
220220
}
221221

222+
/** Identifies an implicit empty first version frozen by handoff, never an explicit write or restore. */
223+
function isMaterializedEmptyInitialVersion(
224+
head: WorkspaceFileVersionSummaryRow | undefined,
225+
file: WorkspaceFileRow
226+
): head is WorkspaceFileVersionSummaryRow {
227+
return (
228+
head !== undefined &&
229+
head.fileId === file.id &&
230+
head.version === INITIAL_WORKSPACE_FILE_VERSION &&
231+
isVersionHeadCurrent(head, file) &&
232+
getWorkspaceFileSize(file) === 0 &&
233+
head.sizeBytes === 0 &&
234+
head.contentHash === null &&
235+
(head.source === 'upload' || head.source === 'unknown') &&
236+
head.restoredFromVersion === null &&
237+
head.createdAt.getTime() === file.contentUpdatedAt.getTime() &&
238+
head.updatedAt.getTime() === file.contentUpdatedAt.getTime()
239+
)
240+
}
241+
222242
/**
223243
* Records a committed content write in the file's history. Runs inside the content-write
224244
* transaction, under the file row's lock, which serializes version numbering and coalescing
225245
* decisions per file.
226246
*
227247
* An empty file with no history is a shell whose content arrives in this write (a create followed
228-
* by its first content), so the shell is not kept as a version of its own.
248+
* by its first content), so the shell is not kept as a version of its own. A handoff may have
249+
* materialized that implicit shell to freeze attribution; its first write still replaces version 1.
229250
*/
230251
export async function recordWorkspaceFileVersionInTx(
231252
tx: DbTransaction,
@@ -239,6 +260,29 @@ export async function recordWorkspaceFileVersionInTx(
239260
return { version: 1, releasedKeys: [previous.key] }
240261
}
241262

263+
if (isMaterializedEmptyInitialVersion(head, previous)) {
264+
await tx
265+
.update(workspaceFileVersion)
266+
.set({
267+
...contentColumns(next, params.nextProvenance),
268+
contentHash: params.contentHash,
269+
source: write.source,
270+
authorUserIds: write.authorUserId ? [write.authorUserId] : [],
271+
restoredFromVersion: write.restoredFromVersion ?? null,
272+
createdAt: now,
273+
updatedAt: now,
274+
})
275+
.where(eq(workspaceFileVersion.id, head.id))
276+
const [references] = await tx.execute<{ referenced: boolean }>(sql`
277+
SELECT EXISTS(SELECT 1 FROM ${workspaceFiles} WHERE ${workspaceFiles.key} = ${previous.key})
278+
OR EXISTS(SELECT 1 FROM ${workspaceFileVersion} WHERE ${workspaceFileVersion.key} = ${previous.key}) AS referenced
279+
`)
280+
return {
281+
version: INITIAL_WORKSPACE_FILE_VERSION,
282+
releasedKeys: references.referenced ? [] : [previous.key],
283+
}
284+
}
285+
242286
if (!head || !isVersionHeadCurrent(head, previous)) {
243287
if (!params.previousProvenance) {
244288
throw new Error('Outgoing workspace file content needs a provenance snapshot to be versioned')

0 commit comments

Comments
 (0)