Skip to content

Commit 3174e38

Browse files
committed
Merge remote-tracking branch 'origin/staging' into codex/durable-agent-memory
2 parents 40cf7d5 + d9e8c4a commit 3174e38

15 files changed

Lines changed: 453 additions & 84 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -868,7 +868,7 @@
868868
"get": {
869869
"operationId": "listFileVersions",
870870
"summary": "List File Versions",
871-
"description": "List the versions of a file, newest first by default. Each write that changes the bytes records one; identical rewrites do not. Collaborative edits within ten minutes fold into one version, as do repeated workflow writes by one author. Renames and moves are not versions. An empty file is version 1 until its first content replaces it. Retention keeps the newest ten and removes older versions by plan, so numbers can have gaps.\n\nOAuth scope: `api:read`.",
871+
"description": "List the versions of a file, newest first by default. Each write that changes the bytes records one; identical rewrites do not. Collaborative edits, and repeated workflow writes by one author, fold into a version under ten minutes old and written in the last five. Renames and moves are not versions. Retention removes older versions by age and plan but keeps the newest ten, so numbers can have gaps.\n\nOAuth scope: `api:read`.",
872872
"x-sim-operation": "files.versions.list",
873873
"x-oauth-scope": "api:read",
874874
"tags": ["Files"],

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

Lines changed: 30 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import { task } from '@trigger.dev/sdk'
66
import { and, count, gt, inArray, isNotNull, lt, min, or, sql } from 'drizzle-orm'
77
import type { CleanupJobPayload } from '@/lib/billing/cleanup-dispatcher'
88
import {
9+
DEFAULT_BATCH_SIZE,
910
DEFAULT_DELETE_CHUNK_SIZE,
1011
DEFAULT_MAX_BATCHES_PER_TABLE,
1112
DEFAULT_WORKSPACE_CHUNK_SIZE,
@@ -22,6 +23,12 @@ const cleanupDb = dbFor('cleanup')
2223
/** Candidate files whose histories are ranked in one query. */
2324
const FILES_PER_QUERY = 500
2425

26+
/**
27+
* Bounds one run like the other cleanup jobs: {@link DEFAULT_MAX_BATCHES_PER_TABLE} batches per
28+
* workspace chunk and this many versions overall. The next run resumes where this one stopped.
29+
*/
30+
const MAX_VERSIONS_PER_RUN = DEFAULT_BATCH_SIZE * DEFAULT_MAX_BATCHES_PER_TABLE
31+
2532
/**
2633
* Superseded versions a free file keeps (its newest 100 with the current one); versions beyond it
2734
* are pruned whatever their age. Paid plans are bounded only by the inline write-time ceiling.
@@ -73,7 +80,12 @@ async function selectCandidateFileIds(
7380
* Superseded versions of the given files past retention: older than the cutoff or beyond the plan's
7481
* count, but never among the newest {@link KEEP_SUPERSEDED} superseded versions of a file.
7582
*/
76-
function selectExpiredVersions(fileIds: string[], cutoff: Date, maxSuperseded: number) {
83+
function selectExpiredVersions(
84+
fileIds: string[],
85+
cutoff: Date,
86+
maxSuperseded: number,
87+
batchSize: number
88+
) {
7789
const ranked = cleanupDb
7890
.select({
7991
id: workspaceFileVersion.id,
@@ -100,7 +112,7 @@ function selectExpiredVersions(fileIds: string[], cutoff: Date, maxSuperseded: n
100112
or(lt(ranked.supersededAt, cutoff), gt(ranked.rank, maxSuperseded))
101113
)
102114
)
103-
.limit(DEFAULT_DELETE_CHUNK_SIZE)
115+
.limit(batchSize)
104116
}
105117

106118
/**
@@ -146,16 +158,27 @@ export async function runCleanupFileVersions(payload: CleanupJobPayload): Promis
146158
)
147159

148160
let deleted = 0
161+
let attempted = 0
149162
for (const group of chunkArray(workspaceIds, DEFAULT_WORKSPACE_CHUNK_SIZE)) {
163+
if (attempted >= MAX_VERSIONS_PER_RUN) break
150164
const candidates = await selectCandidateFileIds(group, cutoff, maxSuperseded)
165+
let batches = 0
151166
for (const fileIds of chunkArray(candidates, FILES_PER_QUERY)) {
152-
for (let batch = 0; batch < DEFAULT_MAX_BATCHES_PER_TABLE; batch++) {
153-
const expired = await selectExpiredVersions(fileIds, cutoff, maxSuperseded)
154-
if (expired.length === 0) break
155-
const removed = await deleteVersions(expired)
167+
let exhausted = false
168+
while (
169+
!exhausted &&
170+
batches < DEFAULT_MAX_BATCHES_PER_TABLE &&
171+
attempted < MAX_VERSIONS_PER_RUN
172+
) {
173+
batches++
174+
const batchSize = Math.min(DEFAULT_DELETE_CHUNK_SIZE, MAX_VERSIONS_PER_RUN - attempted)
175+
const expired = await selectExpiredVersions(fileIds, cutoff, maxSuperseded, batchSize)
176+
attempted += expired.length
177+
const removed = expired.length > 0 ? await deleteVersions(expired) : 0
156178
deleted += removed
157-
if (expired.length < DEFAULT_DELETE_CHUNK_SIZE || removed === 0) break
179+
exhausted = expired.length < batchSize || removed === 0
158180
}
181+
if (!exhausted) break
159182
}
160183
}
161184

‎apps/sim/lib/api/contracts/v2/openapi/files-audit.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -447,7 +447,7 @@ const declaredRoutes = [
447447
operationId: 'listFileVersions',
448448
summary: 'List File Versions',
449449
description:
450-
'List the versions of a file, newest first by default. Each write that changes the bytes records one; identical rewrites do not. Collaborative edits within ten minutes fold into one version, as do repeated workflow writes by one author. Renames and moves are not versions. An empty file is version 1 until its first content replaces it. Retention keeps the newest ten and removes older versions by plan, so numbers can have gaps.',
450+
'List the versions of a file, newest first by default. Each write that changes the bytes records one; identical rewrites do not. Collaborative edits, and repeated workflow writes by one author, fold into a version under ten minutes old and written in the last five. Renames and moves are not versions. Retention removes older versions by age and plan but keeps the newest ten, so numbers can have gaps.',
451451
errors: RESOURCE_ERRORS,
452452
success: { description: 'A page of file versions.' },
453453
}),

‎apps/sim/lib/api/mcp/generated/v2-operations.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1425,7 +1425,7 @@ export const V2_MCP_OPERATIONS = {
14251425
contract: v2ListFileVersionsContract,
14261426
summary: 'List File Versions',
14271427
description:
1428-
'List the versions of a file, newest first by default. Each write that changes the bytes records one; identical rewrites do not. Collaborative edits within ten minutes fold into one version, as do repeated workflow writes by one author. Renames and moves are not versions. An empty file is version 1 until its first content replaces it. Retention keeps the newest ten and removes older versions by plan, so numbers can have gaps.\n\nOAuth scope: `api:read`.',
1428+
'List the versions of a file, newest first by default. Each write that changes the bytes records one; identical rewrites do not. Collaborative edits, and repeated workflow writes by one author, fold into a version under ten minutes old and written in the last five. Renames and moves are not versions. Retention removes older versions by age and plan but keeps the newest ten, so numbers can have gaps.\n\nOAuth scope: `api:read`.',
14291429
handler: () => import('@/app/api/v2/files/[fileId]/versions/route').then((route) => route.GET),
14301430
},
14311431
listKnowledgeBases: {

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

Lines changed: 111 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
/** Real PostgreSQL transactions and local object storage for workspace file version history. */
22
import { mkdtempSync } from 'node:fs'
3-
import { access, rm } from 'node:fs/promises'
3+
import { access, mkdir, rm, writeFile } from 'node:fs/promises'
44
import { tmpdir } from 'node:os'
55
import path from 'node:path'
66
import { db, dbFor } from '@sim/db'
@@ -39,6 +39,7 @@ import {
3939
import { WORKSPACE_FILE_STORAGE_CLEANUP_OUTBOX_EVENT } from '@/lib/uploads/contexts/workspace/workspace-file-storage-cleanup-outbox'
4040
import {
4141
getCurrentWorkspaceFileVersion,
42+
getWorkspaceFileVersion,
4243
queryWorkspaceFileVersions,
4344
releaseWorkspaceFileVersionsForPurgeInTx,
4445
} from '@/lib/uploads/contexts/workspace/workspace-file-versions'
@@ -288,16 +289,123 @@ describe('workspace file version history in PostgreSQL', () => {
288289
)
289290

290291
await expect(deleteWorkspaceFileVersion(fixture.workspaceId, fixture.fileId, 2)).resolves.toBe(
291-
false
292+
'newest'
292293
)
293294
await expect(deleteWorkspaceFileVersion(fixture.workspaceId, fixture.fileId, 1)).resolves.toBe(
294-
true
295+
'deleted'
296+
)
297+
await expect(deleteWorkspaceFileVersion(fixture.workspaceId, fixture.fileId, 1)).resolves.toBe(
298+
'not_found'
295299
)
296300

297301
expect((await versionRows(fixture.fileId)).map((row) => row.version)).toEqual([2])
298302
expect(await objectExists(fixture.firstKey)).toBe(false)
299303
})
300304

305+
it('reads bytes a write left unrecorded as the version its next write records them as', async () => {
306+
const fixture = await seedFile('original')
307+
const write = { source: 'api', authorUserId: fixture.aliceId } as const
308+
await updateWorkspaceFileContent(
309+
fixture.workspaceId,
310+
fixture.fileId,
311+
fixture.aliceId,
312+
Buffer.from('second'),
313+
undefined,
314+
{ version: write }
315+
)
316+
/** A content write that replaced the bytes without recording a version, as a build predating history would. */
317+
const unrecordedKey = `${fixture.firstKey}-unrecorded`
318+
const unrecordedContent = 'third, never recorded'
319+
const unrecordedPath = path.join(fixtureStorage.root, unrecordedKey)
320+
await mkdir(path.dirname(unrecordedPath), { recursive: true })
321+
await writeFile(unrecordedPath, unrecordedContent)
322+
await db
323+
.update(workspaceFiles)
324+
.set({
325+
key: unrecordedKey,
326+
sizeBytes: Buffer.byteLength(unrecordedContent),
327+
contentUpdatedAt: new Date(),
328+
})
329+
.where(eq(workspaceFiles.id, fixture.fileId))
330+
const file = await getWorkspaceFile(fixture.workspaceId, fixture.fileId)
331+
if (!file) throw new Error('file missing')
332+
333+
await expect(
334+
getWorkspaceFileWithCurrentVersion(fixture.workspaceId, fixture.fileId)
335+
).resolves.toMatchObject({ key: unrecordedKey, currentVersion: 3 })
336+
expect(await getCurrentWorkspaceFileVersion(file)).toMatchObject({
337+
version: 3,
338+
key: unrecordedKey,
339+
isCurrent: true,
340+
})
341+
const listed = await queryWorkspaceFileVersions(file, { sortOrder: 'desc', limit: 10 })
342+
expect(listed.versions[0].size).toBe(Buffer.byteLength(unrecordedContent))
343+
expect(listed.versions.map((row) => [row.version, row.isCurrent, row.source])).toEqual([
344+
[3, true, 'unknown'],
345+
[2, false, 'api'],
346+
[1, false, 'upload'],
347+
])
348+
expect(listed.versions[1].supersededAt).toEqual(file.contentUpdatedAt)
349+
350+
const firstPage = await queryWorkspaceFileVersions(file, { sortOrder: 'asc', limit: 2 })
351+
expect(firstPage.versions.map((row) => row.version)).toEqual([1, 2])
352+
const lastPage = await queryWorkspaceFileVersions(file, {
353+
sortOrder: 'asc',
354+
limit: 2,
355+
after: firstPage.nextKeys ?? undefined,
356+
})
357+
expect(lastPage.versions.map((row) => row.version)).toEqual([3])
358+
expect(lastPage.nextKeys).toBeNull()
359+
360+
await expect(deleteWorkspaceFileVersion(fixture.workspaceId, fixture.fileId, 2)).resolves.toBe(
361+
'newest'
362+
)
363+
364+
const next = await updateWorkspaceFileContent(
365+
fixture.workspaceId,
366+
fixture.fileId,
367+
fixture.aliceId,
368+
Buffer.from('fourth'),
369+
undefined,
370+
{ version: write }
371+
)
372+
expect(next.currentVersion).toBe(4)
373+
expect((await versionRows(fixture.fileId)).map((row) => [row.version, row.source])).toEqual([
374+
[1, 'upload'],
375+
[2, 'api'],
376+
[3, 'unknown'],
377+
[4, 'api'],
378+
])
379+
const materialized = (await versionRows(fixture.fileId))[2]
380+
expect(materialized.key).toBe(unrecordedKey)
381+
expect(materialized.sizeBytes).toBe(Buffer.byteLength(unrecordedContent))
382+
expect(await readVersionBytes(fixture.workspaceId, fixture.fileId, materialized.key)).toBe(
383+
unrecordedContent
384+
)
385+
})
386+
387+
it('reads versions against the file as committed, not a record loaded before a write', async () => {
388+
const fixture = await seedFile('original')
389+
const stale = await getWorkspaceFile(fixture.workspaceId, fixture.fileId)
390+
if (!stale) throw new Error('file missing')
391+
await updateWorkspaceFileContent(
392+
fixture.workspaceId,
393+
fixture.fileId,
394+
fixture.aliceId,
395+
Buffer.from('second'),
396+
undefined,
397+
{ version: { source: 'api', authorUserId: fixture.aliceId } }
398+
)
399+
400+
const listed = await queryWorkspaceFileVersions(stale, { sortOrder: 'desc', limit: 10 })
401+
expect(listed.versions.map((row) => [row.version, row.isCurrent])).toEqual([
402+
[2, true],
403+
[1, false],
404+
])
405+
expect((await getCurrentWorkspaceFileVersion(stale)).version).toBe(2)
406+
await expect(getWorkspaceFileVersion(stale, 3)).resolves.toBeNull()
407+
})
408+
301409
it('never folds deliberate writes, and repoints the head for identical bytes', async () => {
302410
const fixture = await seedFile('original')
303411
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
}

0 commit comments

Comments
 (0)