Skip to content

Commit a057433

Browse files
committed
fix(files): agree on the current version when a write skipped recording, and tidy CLI version output
1 parent 339f192 commit a057433

11 files changed

Lines changed: 313 additions & 60 deletions

File tree

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

Lines changed: 73 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -288,16 +288,87 @@ describe('workspace file version history in PostgreSQL', () => {
288288
)
289289

290290
await expect(deleteWorkspaceFileVersion(fixture.workspaceId, fixture.fileId, 2)).resolves.toBe(
291-
false
291+
'newest'
292292
)
293293
await expect(deleteWorkspaceFileVersion(fixture.workspaceId, fixture.fileId, 1)).resolves.toBe(
294-
true
294+
'deleted'
295+
)
296+
await expect(deleteWorkspaceFileVersion(fixture.workspaceId, fixture.fileId, 1)).resolves.toBe(
297+
'not_found'
295298
)
296299

297300
expect((await versionRows(fixture.fileId)).map((row) => row.version)).toEqual([2])
298301
expect(await objectExists(fixture.firstKey)).toBe(false)
299302
})
300303

304+
it('reads bytes a write left unrecorded as the version its next write records them as', async () => {
305+
const fixture = await seedFile('original')
306+
const write = { source: 'api', authorUserId: fixture.aliceId } as const
307+
await updateWorkspaceFileContent(
308+
fixture.workspaceId,
309+
fixture.fileId,
310+
fixture.aliceId,
311+
Buffer.from('second'),
312+
undefined,
313+
{ version: write }
314+
)
315+
/** A content write that replaced the bytes without recording a version, as a build predating history would. */
316+
const unrecordedKey = `${fixture.firstKey}-unrecorded`
317+
await db
318+
.update(workspaceFiles)
319+
.set({ key: unrecordedKey, contentUpdatedAt: new Date() })
320+
.where(eq(workspaceFiles.id, fixture.fileId))
321+
const file = await getWorkspaceFile(fixture.workspaceId, fixture.fileId)
322+
if (!file) throw new Error('file missing')
323+
324+
await expect(
325+
getWorkspaceFileWithCurrentVersion(fixture.workspaceId, fixture.fileId)
326+
).resolves.toMatchObject({ key: unrecordedKey, currentVersion: 3 })
327+
expect(await getCurrentWorkspaceFileVersion(file)).toMatchObject({
328+
version: 3,
329+
key: unrecordedKey,
330+
isCurrent: true,
331+
})
332+
const listed = await queryWorkspaceFileVersions(file, { sortOrder: 'desc', limit: 10 })
333+
expect(listed.versions.map((row) => [row.version, row.isCurrent, row.source])).toEqual([
334+
[3, true, 'unknown'],
335+
[2, false, 'api'],
336+
[1, false, 'upload'],
337+
])
338+
expect(listed.versions[1].supersededAt).toEqual(file.contentUpdatedAt)
339+
340+
const firstPage = await queryWorkspaceFileVersions(file, { sortOrder: 'asc', limit: 2 })
341+
expect(firstPage.versions.map((row) => row.version)).toEqual([1, 2])
342+
const lastPage = await queryWorkspaceFileVersions(file, {
343+
sortOrder: 'asc',
344+
limit: 2,
345+
after: firstPage.nextKeys ?? undefined,
346+
})
347+
expect(lastPage.versions.map((row) => row.version)).toEqual([3])
348+
expect(lastPage.nextKeys).toBeNull()
349+
350+
await expect(deleteWorkspaceFileVersion(fixture.workspaceId, fixture.fileId, 2)).resolves.toBe(
351+
'newest'
352+
)
353+
354+
const next = await updateWorkspaceFileContent(
355+
fixture.workspaceId,
356+
fixture.fileId,
357+
fixture.aliceId,
358+
Buffer.from('fourth'),
359+
undefined,
360+
{ version: write }
361+
)
362+
expect(next.currentVersion).toBe(4)
363+
expect((await versionRows(fixture.fileId)).map((row) => [row.version, row.source])).toEqual([
364+
[1, 'upload'],
365+
[2, 'api'],
366+
[3, 'unknown'],
367+
[4, 'api'],
368+
])
369+
expect((await versionRows(fixture.fileId))[2].key).toBe(unrecordedKey)
370+
})
371+
301372
it('never folds deliberate writes, and repoints the head for identical bytes', async () => {
302373
const fixture = await seedFile('original')
303374
const write = { source: 'api', authorUserId: fixture.aliceId } as const

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

Lines changed: 12 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,7 @@ import {
8282
listWorkspaceFileVersionKeysInTx,
8383
loadWorkspaceFileVersionHead,
8484
recordWorkspaceFileVersionInTx,
85+
type WorkspaceFileVersionDeletion,
8586
type WorkspaceFileVersionWrite,
8687
} from '@/lib/uploads/contexts/workspace/workspace-file-versions'
8788
import { buildStorageKeySegment } from '@/lib/uploads/core/storage-key'
@@ -2097,33 +2098,35 @@ export async function updateWorkspaceFileContent(
20972098
/**
20982099
* Deletes one superseded version of an active workspace file and releases its stored object. The
20992100
* file row is locked so the delete serializes with content writes that supersede or prune history.
2100-
* Returns false when the file has no such superseded version (it never existed, retention removed
2101-
* it, or it is the current version).
21022101
*/
21032102
export async function deleteWorkspaceFileVersion(
21042103
workspaceId: string,
21052104
fileId: string,
21062105
version: number
2107-
): Promise<boolean> {
2108-
const cleanupEventIds = await db.transaction(async (tx) => {
2106+
): Promise<WorkspaceFileVersionDeletion['status']> {
2107+
const deletion = await db.transaction(async (tx) => {
21092108
const [file] = await tx
21102109
.select({ id: workspaceFiles.id })
21112110
.from(workspaceFiles)
21122111
.where(and(eq(workspaceFiles.id, fileId), workspaceFileScopeCondition(workspaceId, 'active')))
21132112
.for('update')
21142113
.limit(1)
21152114
if (!file) throw new OrchestrationError('not_found', 'File not found')
2116-
const key = await deleteWorkspaceFileVersionInTx(tx, fileId, version)
2117-
return key ? enqueueWorkspaceFileStorageCleanups(tx, [key]) : null
2115+
const result = await deleteWorkspaceFileVersionInTx(tx, fileId, version)
2116+
return result.status === 'deleted'
2117+
? {
2118+
status: result.status,
2119+
cleanupEventIds: await enqueueWorkspaceFileStorageCleanups(tx, [result.key]),
2120+
}
2121+
: { status: result.status, cleanupEventIds: [] }
21182122
})
2119-
if (!cleanupEventIds) return false
21202123

2121-
await processWorkspaceFileStorageCleanupsNow(cleanupEventIds, {
2124+
await processWorkspaceFileStorageCleanupsNow(deletion.cleanupEventIds, {
21222125
workspaceId,
21232126
fileId,
21242127
reason: 'deleted version',
21252128
})
2126-
return true
2129+
return deletion.status
21272130
}
21282131

21292132
/**

‎apps/sim/lib/uploads/contexts/workspace/workspace-file-storage-cleanup-outbox.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import type { db } from '@sim/db'
22
import { createLogger } from '@sim/logger'
3-
import { describeError, getErrorMessage } from '@sim/utils/errors'
3+
import { describeError } from '@sim/utils/errors'
44
import { chunkArray } from '@sim/utils/helpers'
55
import {
66
enqueueOutboxEvents,
@@ -87,7 +87,7 @@ export async function processWorkspaceFileStorageCleanupsNow(
8787
logger.warn('Storage cleanup deferred after inline processing error', {
8888
...logContext,
8989
eventId,
90-
error: getErrorMessage(error),
90+
error: describeError(error),
9191
})
9292
}
9393
}

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

Lines changed: 102 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -286,15 +286,22 @@ async function insertVersion(
286286
})
287287
}
288288

289+
/** The outcome of deleting one version; a deleted version's key is released after commit. */
290+
export type WorkspaceFileVersionDeletion =
291+
| { status: 'deleted'; key: string }
292+
| { status: 'not_found' }
293+
| { status: 'newest' }
294+
289295
/**
290-
* Deletes one superseded version, returning its storage key for release after commit, or null when
291-
* no such superseded version exists. The current version is never deletable: it is the file.
296+
* Deletes one superseded version. The newest row is never deleted: it is either the current version
297+
* or the last one recorded before bytes a later write has not recorded yet, and removing it would let
298+
* the next write reuse its number.
292299
*/
293300
export async function deleteWorkspaceFileVersionInTx(
294301
tx: DbTransaction,
295302
fileId: string,
296303
version: number
297-
): Promise<string | null> {
304+
): Promise<WorkspaceFileVersionDeletion> {
298305
const [deleted] = await tx
299306
.delete(workspaceFileVersion)
300307
.where(
@@ -305,7 +312,13 @@ export async function deleteWorkspaceFileVersionInTx(
305312
)
306313
)
307314
.returning({ key: workspaceFileVersion.key })
308-
return deleted?.key ?? null
315+
if (deleted) return { status: 'deleted', key: deleted.key }
316+
const [kept] = await tx
317+
.select({ version: workspaceFileVersion.version })
318+
.from(workspaceFileVersion)
319+
.where(and(eq(workspaceFileVersion.fileId, fileId), eq(workspaceFileVersion.version, version)))
320+
.limit(1)
321+
return kept ? { status: 'newest' } : { status: 'not_found' }
309322
}
310323

311324
/**
@@ -411,7 +424,16 @@ function toSnapshotStatus(status: string | null): WorkspaceFileSecretProvenanceS
411424
return 'unknown'
412425
}
413426

414-
function toVersionRecord(row: WorkspaceFileVersionSummaryRow): WorkspaceFileVersionRecord {
427+
/**
428+
* A stored version as readers see it. A row is current only while it still holds the file's bytes;
429+
* a newest row a later write has replaced without recording (see {@link implicitCurrentVersion})
430+
* reads as superseded from the moment that write landed.
431+
*/
432+
function toVersionRecord(
433+
row: WorkspaceFileVersionSummaryRow,
434+
file: WorkspaceFileVersionSubject
435+
): WorkspaceFileVersionRecord {
436+
const isCurrent = row.supersededAt === null && row.key === file.key
415437
return {
416438
fileId: row.fileId,
417439
version: row.version,
@@ -421,24 +443,35 @@ function toVersionRecord(row: WorkspaceFileVersionSummaryRow): WorkspaceFileVers
421443
source: row.source,
422444
authorUserIds: row.authorUserIds,
423445
restoredFromVersion: row.restoredFromVersion,
424-
isCurrent: row.supersededAt === null,
446+
isCurrent,
425447
createdAt: row.createdAt,
426448
updatedAt: row.updatedAt,
427-
supersededAt: row.supersededAt,
449+
supersededAt: isCurrent ? null : (row.supersededAt ?? contentVersionTime(file)),
428450
}
429451
}
430452

453+
function contentVersionTime(file: WorkspaceFileVersionSubject): Date {
454+
return file.contentUpdatedAt ?? file.updatedAt
455+
}
456+
431457
/**
432-
* Version 1 of a file with no history yet — every file before its first content write. Attributed
433-
* exactly as {@link recordWorkspaceFileVersionInTx} will materialize it, so the version a reader sees
434-
* keeps its identity once it becomes a row.
458+
* The version a file's current bytes hold while no row records them, for a caller that has found
459+
* {@link isVersionHeadCurrent} false. That is version 1 of a file with no history, or the number
460+
* after a newest row that describes other bytes — left by a content write that skipped recording,
461+
* such as one from a build that predates version history. Numbered and attributed exactly as
462+
* {@link recordWorkspaceFileVersionInTx} will materialize it on the next write, so the version a
463+
* reader sees keeps its identity once it becomes a row.
435464
*/
436-
function implicitFirstVersion(file: WorkspaceFileVersionSubject): WorkspaceFileVersionRecord {
437-
const contentUpdatedAt = file.contentUpdatedAt ?? file.updatedAt
438-
const original = isOriginalUploadContent({ uploadedAt: file.uploadedAt, contentUpdatedAt })
465+
function implicitCurrentVersion(
466+
file: WorkspaceFileVersionSubject,
467+
head: WorkspaceFileVersionSummaryRow | undefined
468+
): WorkspaceFileVersionRecord {
469+
const contentUpdatedAt = contentVersionTime(file)
470+
const original =
471+
!head && isOriginalUploadContent({ uploadedAt: file.uploadedAt, contentUpdatedAt })
439472
return {
440473
fileId: file.id,
441-
version: 1,
474+
version: (head?.version ?? 0) + 1,
442475
key: file.key,
443476
size: file.size,
444477
contentType: file.type,
@@ -462,19 +495,32 @@ export async function queryWorkspaceFileVersions(
462495
options: { sortOrder: ListSortOrder; limit: number; after?: CursorKey[] }
463496
): Promise<{ versions: WorkspaceFileVersionRecord[]; nextKeys: CursorKey[] | null }> {
464497
const resume = resumeKeyset(VERSION_KEYSET, options.after, options.sortOrder)
465-
const rows = await db
466-
.select(versionSummaryColumns)
467-
.from(workspaceFileVersion)
468-
.where(and(eq(workspaceFileVersion.fileId, file.id), resume))
469-
.orderBy(...listOrderBy(keysetColumns(VERSION_KEYSET), options.sortOrder))
470-
.limit(options.limit + 1)
471-
const records = rows.map(toVersionRecord)
498+
const [rows, head] = await Promise.all([
499+
db
500+
.select(versionSummaryColumns)
501+
.from(workspaceFileVersion)
502+
.where(and(eq(workspaceFileVersion.fileId, file.id), resume))
503+
.orderBy(...listOrderBy(keysetColumns(VERSION_KEYSET), options.sortOrder))
504+
.limit(options.limit + 1),
505+
loadWorkspaceFileVersionHead(file.id),
506+
])
507+
const records = rows.map((row) => toVersionRecord(row, file))
472508
/**
473-
* An uncursored page that is empty means the file has no rows, so it lists its implicit version 1.
474-
* Any cursor was minted past that version already, so only the first page can carry it.
509+
* The implicit current version numbers above every row, so it leads a descending list and ends an
510+
* ascending one; a cursor already past it leaves it out. Over-fetching by one row still decides
511+
* whether another page follows, since the cut keeps the first `limit` records either way.
475512
*/
476-
if (records.length === 0 && !options.after) {
477-
records.push(implicitFirstVersion(file))
513+
const implicit = isVersionHeadCurrent(head, file) ? null : implicitCurrentVersion(file, head)
514+
const resumeAfter = options.after?.[0]
515+
if (
516+
implicit &&
517+
(typeof resumeAfter !== 'number' ||
518+
(options.sortOrder === 'desc'
519+
? implicit.version < resumeAfter
520+
: implicit.version > resumeAfter))
521+
) {
522+
if (options.sortOrder === 'desc') records.unshift(implicit)
523+
else records.push(implicit)
478524
}
479525
const page = keysetPage(VERSION_KEYSET, records, options.limit)
480526
return { versions: page.data, nextKeys: page.nextCursorKeys }
@@ -495,28 +541,44 @@ export async function findWorkspaceFileVersionKeys(keys: readonly string[]): Pro
495541
return new Set(rows.map((row) => row.key))
496542
}
497543

498-
/** The current version, or the implicit version 1 of a file with no history rows. */
544+
/** The version holding the file's current bytes, recorded or implicit. */
499545
export async function getCurrentWorkspaceFileVersion(
500546
file: WorkspaceFileVersionSubject
501547
): Promise<WorkspaceFileVersionRecord> {
502548
const head = await loadWorkspaceFileVersionHead(file.id)
503-
return head ? toVersionRecord(head) : implicitFirstVersion(file)
549+
return head && isVersionHeadCurrent(head, file)
550+
? toVersionRecord(head, file)
551+
: implicitCurrentVersion(file, head)
504552
}
505553

506554
/**
507555
* The current version number of the enclosing query's `workspace_files` row, as a correlated
508-
* subquery so the row and its number come from one statement's snapshot. Content writes record
509-
* their version in the transaction that replaces the bytes, so the newest row describes the current
510-
* content; a file with no rows is on its implicit version 1. Both sides of the correlation
511-
* are table-qualified because Drizzle renders single-table columns bare, which would bind the outer
512-
* `id` to this subquery's own table.
556+
* subquery so the row and its number come from one statement's snapshot. The newest row numbers
557+
* the file's bytes while it still holds them; otherwise the bytes are on the implicit version after
558+
* it, and a file with no rows is on version 1 — the numbering {@link implicitCurrentVersion} gives.
559+
* Both sides of the correlation are table-qualified because Drizzle renders single-table columns
560+
* bare, which would bind the outer columns to this subquery's own table.
513561
*/
514562
export function currentWorkspaceFileVersionNumberSql() {
515-
const versionFileId = sql`${workspaceFileVersion}.${sql.identifier(workspaceFileVersion.fileId.name)}`
516-
const outerFileId = sql`${workspaceFiles}.${sql.identifier(workspaceFiles.id.name)}`
517-
return sql<number>`coalesce((select max(${workspaceFileVersion.version}) from ${workspaceFileVersion} where ${versionFileId} = ${outerFileId}), 1)`.mapWith(
518-
Number
519-
)
563+
const qualified = (table: typeof workspaceFileVersion | typeof workspaceFiles, name: string) =>
564+
sql`${table}.${sql.identifier(name)}`
565+
const head = {
566+
fileId: qualified(workspaceFileVersion, workspaceFileVersion.fileId.name),
567+
version: qualified(workspaceFileVersion, workspaceFileVersion.version.name),
568+
key: qualified(workspaceFileVersion, workspaceFileVersion.key.name),
569+
supersededAt: qualified(workspaceFileVersion, workspaceFileVersion.supersededAt.name),
570+
}
571+
const file = {
572+
id: qualified(workspaceFiles, workspaceFiles.id.name),
573+
key: qualified(workspaceFiles, workspaceFiles.key.name),
574+
}
575+
const number = sql`case when ${head.supersededAt} is null and ${head.key} = ${file.key} then ${head.version} else ${head.version} + 1 end`
576+
return sql<number>`coalesce((
577+
select ${number} from ${workspaceFileVersion}
578+
where ${head.fileId} = ${file.id}
579+
order by ${head.version} desc
580+
limit 1
581+
), 1)`.mapWith(Number)
520582
}
521583

522584
/** One version of a file, or null when it never existed or retention removed it. */
@@ -529,9 +591,9 @@ export async function getWorkspaceFileVersion(
529591
.from(workspaceFileVersion)
530592
.where(and(eq(workspaceFileVersion.fileId, file.id), eq(workspaceFileVersion.version, version)))
531593
.limit(1)
532-
if (row) return toVersionRecord(row)
533-
if (version !== 1 || (await loadWorkspaceFileVersionHead(file.id))) return null
534-
return implicitFirstVersion(file)
594+
if (row) return toVersionRecord(row, file)
595+
const current = await getCurrentWorkspaceFileVersion(file)
596+
return current.version === version ? current : null
535597
}
536598

537599
/**

0 commit comments

Comments
 (0)