Skip to content

Commit a730e59

Browse files
fix(workspace-fork): scope and batch preview revisions (#7901)
1 parent 6ea53c0 commit a730e59

7 files changed

Lines changed: 352 additions & 23 deletions

File tree

‎.github/workflows/test-build.yml‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,12 @@ jobs:
121121
BILLING_USAGE_TEST_REDIS_URL: redis://127.0.0.1:6379
122122
run: bunx vitest run lib/billing/core/usage-log.postgres.test.ts lib/billing/core/organization-activity.postgres.test.ts lib/billing/calculations/usage-reservation.test.ts
123123

124+
- name: Verify fork previews ignore execution file history in PostgreSQL
125+
working-directory: apps/sim
126+
env:
127+
FORK_REVISION_TEST_DATABASE_URL: postgresql://postgres:postgres@127.0.0.1:5432/sim_auth_scim
128+
run: bunx vitest run ee/workspace-forking/application/revision.postgres.test.ts
129+
124130
- name: Verify cumulative billing timeout recovery on PostgreSQL 16
125131
if: matrix.provision == 'push'
126132
working-directory: apps/sim
Lines changed: 180 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,180 @@
1+
/**
2+
* @vitest-environment node
3+
*
4+
* Set FORK_REVISION_TEST_DATABASE_URL to a local PostgreSQL database. Each run uses an
5+
* isolated schema and executes the real revision query against execution-heavy workspaces.
6+
*/
7+
import * as schema from '@sim/db/schema'
8+
import { generateShortId } from '@sim/utils/id'
9+
import { drizzle } from 'drizzle-orm/postgres-js'
10+
import postgres from 'postgres'
11+
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
12+
import {
13+
assertForkPreviewFresh,
14+
loadForkPreviewRevision,
15+
} from '@/ee/workspace-forking/application/revision'
16+
17+
vi.unmock('drizzle-orm')
18+
vi.unmock('@sim/db/schema')
19+
20+
const databaseUrl = process.env.FORK_REVISION_TEST_DATABASE_URL
21+
if (databaseUrl && !['localhost', '127.0.0.1', '[::1]'].includes(new URL(databaseUrl).hostname)) {
22+
throw new Error('Fork revision PostgreSQL tests require a local database')
23+
}
24+
25+
describe.runIf(Boolean(databaseUrl))('fork revision scope in PostgreSQL', () => {
26+
const testSchema = `fork_revision_${generateShortId()
27+
.replace(/[^a-zA-Z0-9]/g, '')
28+
.toLowerCase()}`
29+
let client: ReturnType<typeof postgres>
30+
let executor: ReturnType<typeof drizzle<typeof schema>>
31+
const scope = {
32+
sourceWorkspaceId: 'source',
33+
targetWorkspaceId: 'target',
34+
edge: { parentWorkspaceId: 'target', childWorkspaceId: 'source' },
35+
}
36+
37+
beforeAll(async () => {
38+
client = postgres(databaseUrl!, { max: 1, connection: { search_path: testSchema } })
39+
executor = drizzle(client, { schema })
40+
await client.unsafe(`CREATE SCHEMA ${testSchema}`)
41+
await client.unsafe(`
42+
CREATE TABLE workspace (
43+
id text PRIMARY KEY, organization_id text, name text,
44+
storage_used_bytes bigint DEFAULT 0, updated_at timestamp
45+
);
46+
CREATE TABLE workflow (
47+
id text PRIMARY KEY, workspace_id text, name text, archived_at timestamp,
48+
fork_sync_excluded boolean DEFAULT false, is_deployed boolean DEFAULT true,
49+
run_count integer DEFAULT 0, last_run_at timestamp, updated_at timestamp
50+
);
51+
CREATE TABLE workflow_deployment_version (
52+
id text PRIMARY KEY, workflow_id text, is_active boolean, state jsonb
53+
);
54+
CREATE TABLE workspace_files (
55+
id text PRIMARY KEY, workspace_id text, context text, deleted_at timestamp,
56+
key text, original_name text, size_bytes bigint, content_updated_at timestamp
57+
);
58+
CREATE INDEX ON workspace_files (workspace_id)
59+
WHERE context = 'workspace' AND deleted_at IS NULL;
60+
CREATE TABLE permissions (id text PRIMARY KEY, entity_id text, entity_type text, permission_type text);
61+
CREATE TABLE custom_block (id text PRIMARY KEY, organization_id text, workflow_id text);
62+
`)
63+
for (const table of ['workflow_blocks', 'workflow_edges', 'workflow_subflows', 'webhook']) {
64+
await client.unsafe(
65+
`CREATE TABLE ${table} (id text PRIMARY KEY, workflow_id text, data jsonb)`
66+
)
67+
}
68+
for (const table of [
69+
'folder',
70+
'user_table_definitions',
71+
'knowledge_base',
72+
'custom_tools',
73+
'skill',
74+
'mcp_servers',
75+
'credential',
76+
'workspace_environment',
77+
'workspace_sandbox',
78+
]) {
79+
await client.unsafe(
80+
`CREATE TABLE ${table} (id text PRIMARY KEY, workspace_id text, data jsonb)`
81+
)
82+
}
83+
for (const table of [
84+
'workspace_fork_resource_map',
85+
'workspace_fork_block_map',
86+
'workspace_fork_dependent_value',
87+
]) {
88+
await client.unsafe(
89+
`CREATE TABLE ${table} (id text PRIMARY KEY, child_workspace_id text, data jsonb)`
90+
)
91+
}
92+
await client`INSERT INTO workspace (id, name) VALUES ('source', 'Source'), ('target', 'Target')`
93+
await client`INSERT INTO workflow (id, workspace_id, name)
94+
VALUES ('source-workflow', 'source', 'Source workflow'), ('target-workflow', 'target', 'Target workflow')`
95+
await client`INSERT INTO workflow_deployment_version (id, workflow_id, is_active, state)
96+
VALUES ('deployment', 'source-workflow', true, '{"blocks":{}}')`
97+
})
98+
99+
afterAll(async () => {
100+
if (!client) return
101+
await client.unsafe(`DROP SCHEMA IF EXISTS ${testSchema} CASCADE`)
102+
await client.end()
103+
})
104+
105+
beforeEach(async () => {
106+
await client`TRUNCATE workspace_files, workflow_blocks, permissions`
107+
await client`UPDATE workspace SET storage_used_bytes = 0, updated_at = null`
108+
await client`INSERT INTO workspace_files
109+
(id, workspace_id, context, key, original_name, size_bytes, content_updated_at)
110+
VALUES ('source-file', 'source', 'workspace', 'workspace/source/file', 'file.txt', 24, '2026-09-01'),
111+
('target-file', 'target', 'workspace', 'workspace/target/file', 'file.txt', 24, '2026-09-01')`
112+
})
113+
114+
it('ignores more than 100,000 execution files, deleted files, chat uploads, and runtime storage changes', async () => {
115+
const before = await loadForkPreviewRevision(executor, scope, {})
116+
await client`INSERT INTO workspace_files (id, workspace_id, context, key, original_name, size_bytes)
117+
SELECT 'execution-' || n, CASE WHEN n % 2 = 0 THEN 'source' ELSE 'target' END,
118+
'execution', 'execution/' || n, repeat('x', 700), 128
119+
FROM generate_series(1, 100001) n`
120+
await client`INSERT INTO workspace_files (id, workspace_id, context, deleted_at)
121+
VALUES ('deleted', 'source', 'workspace', now()), ('upload', 'target', 'mothership', null),
122+
('kb-document', 'source', 'knowledge-base', null), ('other-workspace', 'unrelated', 'workspace', null)`
123+
await client`UPDATE workspace SET storage_used_bytes = 999999, updated_at = now()`
124+
await client`UPDATE workflow SET run_count = run_count + 1, last_run_at = now(), updated_at = now()`
125+
126+
const after = await loadForkPreviewRevision(executor, scope, {})
127+
expect(after).toEqual(before)
128+
await expect(
129+
assertForkPreviewFresh(executor, scope, {
130+
workspaceId: 'source',
131+
requestId: 'request',
132+
requestHash: 'hash',
133+
previewFingerprint: before.fingerprint,
134+
choices: {},
135+
})
136+
).resolves.toBeUndefined()
137+
})
138+
139+
it.each(['source', 'target'])(
140+
'invalidates previews when an active %s file changes or disappears',
141+
async (workspaceId) => {
142+
const before = await loadForkPreviewRevision(executor, scope, {})
143+
await client`UPDATE workspace_files SET content_updated_at = '2026-09-02' WHERE workspace_id = ${workspaceId}`
144+
const edited = await loadForkPreviewRevision(executor, scope, {})
145+
expect(edited.categories.files).not.toBe(before.categories.files)
146+
147+
await client`UPDATE workspace_files SET deleted_at = now() WHERE workspace_id = ${workspaceId}`
148+
const deleted = await loadForkPreviewRevision(executor, scope, {})
149+
expect(deleted.categories.files).not.toBe(edited.categories.files)
150+
151+
await client`UPDATE workspace_files SET deleted_at = null WHERE workspace_id = ${workspaceId}`
152+
expect((await loadForkPreviewRevision(executor, scope, {})).fingerprint).toBe(
153+
edited.fingerprint
154+
)
155+
}
156+
)
157+
158+
it('still detects graph edits, access changes, and changed copy choices', async () => {
159+
const before = await loadForkPreviewRevision(executor, scope, {})
160+
await client`INSERT INTO workflow_blocks (id, workflow_id, data)
161+
VALUES ('block', 'target-workflow', '{"value":"changed"}')`
162+
await client`INSERT INTO permissions (id, entity_id, entity_type, permission_type)
163+
VALUES ('member', 'source', 'workspace', 'admin')`
164+
const after = await loadForkPreviewRevision(executor, scope, {})
165+
expect(after.categories.target_graph).not.toBe(before.categories.target_graph)
166+
expect(after.categories.membership).not.toBe(before.categories.membership)
167+
expect(
168+
(await loadForkPreviewRevision(executor, scope, { copyResources: [] })).fingerprint
169+
).not.toBe(after.fingerprint)
170+
})
171+
172+
it('retains the row limit for actual fork resources', async () => {
173+
await client`INSERT INTO workspace_files (id, workspace_id, context)
174+
SELECT 'durable-' || n, 'source', 'workspace' FROM generate_series(1, 100001) n`
175+
await expect(loadForkPreviewRevision(executor, scope, {})).rejects.toMatchObject({
176+
statusCode: 413,
177+
message: 'Fork preview files exceeds its 100000 row ceiling',
178+
})
179+
})
180+
})
Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
/** @vitest-environment node */
2+
import type { SQL } from 'drizzle-orm'
3+
import { PgDialect } from 'drizzle-orm/pg-core'
4+
import { describe, expect, it, vi } from 'vitest'
5+
import type { DbOrTx } from '@/lib/db/types'
6+
import { WorkspaceOperationConflict } from '@/lib/workspaces/operations/receipts'
7+
import {
8+
assertForkPreviewFresh,
9+
loadForkPreviewRevision,
10+
} from '@/ee/workspace-forking/application/revision'
11+
12+
vi.unmock('drizzle-orm')
13+
vi.unmock('@sim/db/schema')
14+
15+
const scope = {
16+
sourceWorkspaceId: 'source',
17+
targetWorkspaceId: 'target',
18+
edge: { parentWorkspaceId: 'target', childWorkspaceId: 'source' },
19+
}
20+
21+
function mockRevisionExecutor(overrides: { count?: string; bytes?: string; digest?: string } = {}) {
22+
const execute = vi.fn(async (_query: SQL) => [
23+
{ category: 'files', count: '3', bytes: '1024', digest: 'file-revision', ...overrides },
24+
])
25+
return { execute, executor: { execute } as unknown as DbOrTx }
26+
}
27+
28+
describe('fork preview revisions', () => {
29+
it('reads every category in one bounded query and excludes files the sync cannot copy', async () => {
30+
const { execute, executor } = mockRevisionExecutor()
31+
await loadForkPreviewRevision(executor, scope, {})
32+
33+
expect(execute).toHaveBeenCalledTimes(1)
34+
const query = new PgDialect().sqlToQuery(execute.mock.calls[0][0])
35+
expect(query.sql).toMatch(/"workspace_files"\."context" = \$\d+/)
36+
expect(query.sql).toContain('"workspace_files"."deleted_at" is null')
37+
expect(query.params).toContain('workspace')
38+
expect(query.sql).toContain("ARRAY['updated_at', 'storage_used_bytes']")
39+
expect(query.sql.match(/LIMIT \$\d+/g)).toHaveLength(20)
40+
expect(query.params.filter((value) => value === 100_001)).toHaveLength(20)
41+
expect(query.params).toContain('source')
42+
expect(query.params).toContain('target')
43+
expect(query.params).toContain('mappings')
44+
expect(query.params).toContain('block_identities')
45+
expect(query.params).toContain('dependent_values')
46+
expect(query.params.some(Array.isArray)).toBe(false)
47+
})
48+
49+
it('supports creating a fork without a target or existing edge', async () => {
50+
const { execute, executor } = mockRevisionExecutor()
51+
await loadForkPreviewRevision(executor, { sourceWorkspaceId: 'source' }, {})
52+
53+
const query = new PgDialect().sqlToQuery(execute.mock.calls[0][0])
54+
expect(query.sql.match(/LIMIT \$\d+/g)).toHaveLength(17)
55+
expect(query.params).not.toContain('target')
56+
expect(query.params).not.toContain('mappings')
57+
expect(query.params.some((value) => value == null)).toBe(false)
58+
})
59+
60+
it.each([
61+
{ count: '100001', bytes: '1', message: 'Fork preview files exceeds its 100000 row ceiling' },
62+
{
63+
count: '1',
64+
bytes: String(64 * 1024 * 1024 + 1),
65+
message: 'Fork preview files exceeds its 64 MiB byte ceiling',
66+
},
67+
])('rejects an oversized category with the actual limiting budget: $message', async (row) => {
68+
const { executor } = mockRevisionExecutor(row)
69+
await expect(loadForkPreviewRevision(executor, scope, {})).rejects.toMatchObject({
70+
statusCode: 413,
71+
message: row.message,
72+
})
73+
})
74+
75+
it('accepts categories exactly at both limits', async () => {
76+
const { executor } = mockRevisionExecutor({ count: '100000', bytes: String(64 * 1024 * 1024) })
77+
await expect(loadForkPreviewRevision(executor, scope, {})).resolves.toMatchObject({
78+
categories: { files: 'file-revision' },
79+
})
80+
})
81+
82+
it('keeps scope and copy choices bound to the fingerprint', async () => {
83+
const { executor } = mockRevisionExecutor()
84+
const original = await loadForkPreviewRevision(executor, scope, {})
85+
const changedScope = await loadForkPreviewRevision(
86+
executor,
87+
{ ...scope, targetWorkspaceId: 'other-target' },
88+
{}
89+
)
90+
const changedChoices = await loadForkPreviewRevision(executor, scope, { copyResources: [] })
91+
expect(changedScope.fingerprint).not.toBe(original.fingerprint)
92+
expect(changedChoices.fingerprint).not.toBe(original.fingerprint)
93+
})
94+
95+
it('still refuses apply when a reviewed resource changes', async () => {
96+
const { executor, execute } = mockRevisionExecutor()
97+
const preview = await loadForkPreviewRevision(executor, scope, {})
98+
const admission = {
99+
workspaceId: 'source',
100+
requestId: 'request',
101+
requestHash: 'request-hash',
102+
previewFingerprint: preview.fingerprint,
103+
choices: {},
104+
}
105+
await expect(assertForkPreviewFresh(executor, scope, admission)).resolves.toBeUndefined()
106+
execute.mockResolvedValue([{ category: 'files', count: '3', bytes: '1024', digest: 'changed' }])
107+
await expect(assertForkPreviewFresh(executor, scope, admission)).rejects.toBeInstanceOf(
108+
WorkspaceOperationConflict
109+
)
110+
})
111+
})

‎apps/sim/ee/workspace-forking/application/revision.ts‎

Lines changed: 39 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -22,9 +22,10 @@ import {
2222
workspaceForkResourceMap,
2323
workspaceSandbox,
2424
} from '@sim/db/schema'
25-
import { type SQL, sql } from 'drizzle-orm'
25+
import { and, type SQL, sql } from 'drizzle-orm'
2626
import type { DbOrTx } from '@/lib/db/types'
2727
import { acquireFolderMutationLock } from '@/lib/folders/locks'
28+
import { activeWorkspaceFileConditions } from '@/lib/workspace-files/query-scope'
2829
import {
2930
WorkspaceOperationConflict,
3031
workflowOperationFingerprint,
@@ -46,7 +47,14 @@ export interface ForkMutationAdmission {
4647
choices: Record<string, unknown>
4748
}
4849

49-
/** Digests are bounded database aggregates; graph and secret values never enter preview diagnostics. */
50+
const MAX_REVISION_ROWS = 100_000
51+
const MAX_REVISION_BYTES = 64 * 1024 * 1024
52+
53+
/**
54+
* Fingerprints fork configuration in one database snapshot. Runtime file outputs and their
55+
* storage ledger are not sync inputs; including them makes ordinary executions invalidate
56+
* previews. Only bounded aggregates leave the database, never graph or secret values.
57+
*/
5058
export async function loadForkPreviewRevision(
5159
executor: DbOrTx,
5260
scope: ForkRevisionScope,
@@ -64,7 +72,7 @@ export async function loadForkPreviewRevision(
6472
)
6573
const workflowIds = sql`SELECT id FROM ${workflow} WHERE workspace_id IN (${values})`
6674
const queries: Record<string, SQL> = {
67-
workspaces: sql`SELECT id, to_jsonb(r) - ARRAY['updated_at'] AS state FROM ${workspace} r WHERE id IN (${values})`,
75+
workspaces: sql`SELECT id, to_jsonb(r) - ARRAY['updated_at', 'storage_used_bytes'] AS state FROM ${workspace} r WHERE id IN (${values})`,
6876
workflows: sql`SELECT id, to_jsonb(r) - ARRAY['run_count', 'last_run_at', 'last_synced', 'updated_at'] AS state FROM ${workflow} r WHERE workspace_id IN (${values})`,
6977
source_deployments: sql`SELECT d.id, to_jsonb(d) AS state FROM ${workflowDeploymentVersion} d JOIN ${workflow} w ON w.id = d.workflow_id WHERE w.workspace_id = ${scope.sourceWorkspaceId} AND d.is_active = true AND w.archived_at IS NULL AND w.fork_sync_excluded = false`,
7078
target_graph: sql`SELECT 'block:' || b.id AS id, to_jsonb(b) - ARRAY['updated_at', 'created_at'] AS state FROM ${workflowBlocks} b JOIN ${workflow} w ON w.id = b.workflow_id WHERE w.workspace_id = ${scope.targetWorkspaceId ?? scope.sourceWorkspaceId}
@@ -78,7 +86,7 @@ export async function loadForkPreviewRevision(
7886
tools: sql`SELECT id, to_jsonb(r) AS state FROM ${customTools} r WHERE workspace_id IN (${values})`,
7987
skills: sql`SELECT id, to_jsonb(r) AS state FROM ${skill} r WHERE workspace_id IN (${values})`,
8088
servers: sql`SELECT id, to_jsonb(r) - ARRAY['updated_at', 'last_connected_at', 'last_tools_refresh', 'tool_count', 'connection_status', 'last_error'] AS state FROM ${mcpServers} r WHERE workspace_id IN (${values})`,
81-
files: sql`SELECT id, to_jsonb(r) AS state FROM ${workspaceFiles} r WHERE workspace_id IN (${values})`,
89+
files: sql`SELECT id, to_jsonb(${workspaceFiles}) AS state FROM ${workspaceFiles} WHERE ${and(...activeWorkspaceFileConditions(ids))}`,
8290
credentials: sql`SELECT id, to_jsonb(r) - ARRAY['updated_at', 'last_used_at'] AS state FROM ${credential} r WHERE workspace_id IN (${values})`,
8391
secrets: sql`SELECT id, to_jsonb(r) - 'updated_at' AS state FROM ${workspaceEnvironment} r WHERE workspace_id IN (${values})`,
8492
sandboxes: sql`SELECT id, to_jsonb(r) AS state FROM ${workspaceSandbox} r WHERE workspace_id IN (${values})`,
@@ -89,17 +97,34 @@ export async function loadForkPreviewRevision(
8997
queries.block_identities = sql`SELECT id, to_jsonb(r) AS state FROM ${workspaceForkBlockMap} r WHERE child_workspace_id = ${scope.edge.childWorkspaceId}`
9098
queries.dependent_values = sql`SELECT id, to_jsonb(r) AS state FROM ${workspaceForkDependentValue} r WHERE child_workspace_id = ${scope.edge.childWorkspaceId}`
9199
}
92-
const categories: Record<string, string> = {}
93-
for (const [category, rows] of Object.entries(queries)) {
94-
const [size] = await executor.execute<{ count: string; bytes: string }>(
95-
sql`SELECT count(*)::text AS count, coalesce(sum(octet_length(state::text)), 0)::text AS bytes FROM (${rows}) revision_rows`
96-
)
97-
if (Number(size.count) > 100000 || Number(size.bytes) > 64 * 1024 * 1024)
98-
throw new ForkError(`Fork preview ${category} exceeds its row or 64 MiB byte ceiling`, 413)
99-
const [revision] = await executor.execute<{ digest: string }>(
100-
sql`SELECT md5(coalesce(string_agg(md5(state::text), '' ORDER BY id), '')) AS digest FROM (${rows}) revision_rows`
100+
const revisions = await executor.execute<{
101+
category: string
102+
count: string
103+
bytes: string
104+
digest: string
105+
}>(
106+
sql.join(
107+
Object.entries(queries).map(
108+
([category, rows]) => sql`
109+
SELECT ${category}::text AS category, count(*)::text AS count,
110+
coalesce(sum(octet_length(state::text)), 0)::text AS bytes,
111+
md5(coalesce(string_agg(md5(state::text), '' ORDER BY id), '')) AS digest
112+
FROM (SELECT id, state FROM (${rows}) revision_source LIMIT ${MAX_REVISION_ROWS + 1}) revision_rows
113+
`
114+
),
115+
sql` UNION ALL `
101116
)
102-
categories[category] = revision.digest
117+
)
118+
const categories: Record<string, string> = {}
119+
for (const revision of revisions) {
120+
if (Number(revision.count) > MAX_REVISION_ROWS)
121+
throw new ForkError(
122+
`Fork preview ${revision.category} exceeds its ${MAX_REVISION_ROWS} row ceiling`,
123+
413
124+
)
125+
if (Number(revision.bytes) > MAX_REVISION_BYTES)
126+
throw new ForkError(`Fork preview ${revision.category} exceeds its 64 MiB byte ceiling`, 413)
127+
categories[revision.category] = revision.digest
103128
}
104129
return {
105130
categories,

0 commit comments

Comments
 (0)