diff --git a/apps/sim/ee/workspace-forking/application/sync-default.test.ts b/apps/sim/ee/workspace-forking/application/sync-default.test.ts index 5eb0b50df73..accfcd5a4d1 100644 --- a/apps/sim/ee/workspace-forking/application/sync-default.test.ts +++ b/apps/sim/ee/workspace-forking/application/sync-default.test.ts @@ -1,10 +1,9 @@ /** * @vitest-environment node */ -import { workspace } from '@sim/db/schema' import { createSessionPrincipal } from '@sim/testing/factories/principal.factory' -import { auditMock, auditMockFns } from '@sim/testing/mocks/audit.mock' -import { dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing/mocks/database.mock' +import { auditMock } from '@sim/testing/mocks/audit.mock' +import { resetDbChainMock } from '@sim/testing/mocks/database.mock' import { permissionsMock, permissionsMockFns } from '@sim/testing/mocks/permissions.mock' import { posthogServerMock } from '@sim/testing/mocks/posthog-server.mock' import { @@ -33,7 +32,6 @@ vi.mock('@/ee/workspace-forking/lib/lineage/lineage-root', () => workspaceForkin import { setForkSyncDefault } from '@/ee/workspace-forking/application/sync-default' const principal = createSessionPrincipal({ userId: 'actor-1' }) -const LINEAGE = ['root-ws', 'fork-a', 'fork-b', 'grandchild'] beforeEach(() => { resetDbChainMock() @@ -45,50 +43,21 @@ beforeEach(() => { }) workspaceAuthorizationMockFns.mockAuthorizeWorkspaceOperation.mockResolvedValue(undefined) mockResolveForkLineageRootId.mockResolvedValue('root-ws') - mockResolveForkLineageWorkspaceIds.mockResolvedValue(LINEAGE) - // The live-member row lock, then the `UPDATE ... RETURNING` of the members that changed. - queueTableRows( - workspace, - LINEAGE.map((id) => ({ id })) - ) - dbChainMockFns.returning.mockResolvedValue(LINEAGE.map((id) => ({ id, name: `Name of ${id}` }))) }) -const run = (excludeNewWorkflows: boolean) => - setForkSyncDefault.execute({ - principal, - input: { workspaceId: 'fork-a', excludeNewWorkflows }, - }) - +/** + * The lineage-wide write and its audit fan-out are proven against real Postgres in + * `fork-sync.integration.ts`. This covers the one branch that suite cannot stage: an unlink + * committing between the root walk and the lineage lock. + */ describe('setForkSyncDefault', () => { - /** - * One entry per member, because the default genuinely changed for all of them. A single - * entry on the calling workspace would leave the other members' admins with no record. - */ - it("files one audit entry per updated member, in that member's own workspace and name", async () => { - await run(true) - const audited = auditMockFns.mockRecordAudit.mock.calls.map(([entry]) => entry) - expect(audited).toHaveLength(LINEAGE.length) - expect(audited.map((entry) => entry.resourceId).sort()).toEqual([...LINEAGE].sort()) - for (const entry of audited) { - expect(entry.action).toBe('workspace.fork_sync_default_changed') - // Filed in the workspace it describes. Without an explicit workspaceId the wrapper - // defaults it to the caller's workspace, so every entry would pile into one log and - // the other members' admins would see nothing. - expect(entry.workspaceId).toBe(entry.resourceId) - // The member's OWN name, not a raw id and not the caller's name. - expect(entry.resourceName).toBe(`Name of ${entry.resourceId}`) - expect(entry.metadata).toMatchObject({ - forkSyncNewWorkflowsExcluded: true, - originWorkspaceId: 'fork-a', - originWorkspaceName: 'Fork A', - }) - } - }) - - /** An unlink that moved the caller out of the locked root's lineage must not be written. */ it('refuses when the locked root no longer reaches the calling workspace', async () => { mockResolveForkLineageWorkspaceIds.mockResolvedValue(['root-ws', 'fork-b']) - await expect(run(true)).rejects.toMatchObject({ statusCode: 409 }) + await expect( + setForkSyncDefault.execute({ + principal, + input: { workspaceId: 'fork-a', excludeNewWorkflows: true }, + }) + ).rejects.toMatchObject({ statusCode: 409 }) }) }) diff --git a/apps/sim/lib/workspaces/__integration__/fork-sync.integration.ts b/apps/sim/lib/workspaces/__integration__/fork-sync.integration.ts index bb0210d5139..619a4068957 100644 --- a/apps/sim/lib/workspaces/__integration__/fork-sync.integration.ts +++ b/apps/sim/lib/workspaces/__integration__/fork-sync.integration.ts @@ -1,5 +1,7 @@ +import { AuditAction } from '@sim/audit' import { db } from '@sim/db' import { + auditLog, folder, outboxEvent, permissions, @@ -15,8 +17,8 @@ import { workspaceSandbox, } from '@sim/db/schema' import { generateId } from '@sim/utils/id' -import { and, eq } from 'drizzle-orm' -import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { and, eq, inArray, sql } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' import { withWorkspaceInvocationScope } from '@/lib/core/application/workspace-invocation-scope' import { processOutboxEventById } from '@/lib/core/outbox/service' import { createScopedCliTransport } from '@/lib/mothership/agent-cli/scoped-transport' @@ -765,9 +767,10 @@ describe('authorized fork and sync against PostgreSQL', () => { /** * The opt-in policy end to end: set from a fork, it reaches the parent and changes only - * what differs; a genuinely new workflow (created, duplicated, or a fork's starter) and a - * new fork take it; a forked copy stays synced; no existing workflow moves; and an archived - * member is walked through for the lineage root but never written. + * what differs, filing one audit entry in each changed workspace's own log; a genuinely + * new workflow (created, duplicated, or a fork's starter) and a new fork take it; a forked + * copy stays synced; no existing workflow moves; and an archived member is walked through + * for the lineage root but never written. */ it('gives new workflows the lineage fork-sync default while copies stay synced', async () => { const childId = await createChild() @@ -787,6 +790,22 @@ describe('authorized fork and sync against PostgreSQL', () => { )[0]?.excluded const setDefault = (workspaceId: string, excludeNewWorkflows: boolean) => setForkSyncDefault.execute({ principal, input: { workspaceId, excludeNewWorkflows } }) + /** Every audit entry a change issued from `originId` filed, whichever workspace it named. */ + const auditedFrom = (originId: string) => + db + .select({ + workspaceId: auditLog.workspaceId, + resourceId: auditLog.resourceId, + resourceName: auditLog.resourceName, + metadata: auditLog.metadata, + }) + .from(auditLog) + .where( + and( + eq(auditLog.action, AuditAction.WORKSPACE_FORK_SYNC_DEFAULT_CHANGED), + sql`${auditLog.metadata} ->> 'originWorkspaceId' = ${originId}` + ) + ) try { const first = await setDefault(childId, true) expect(first.changedWorkspaces.map((member) => member.id)).toEqual( @@ -794,6 +813,41 @@ describe('authorized fork and sync against PostgreSQL', () => { ) expect((await setDefault(childId, true)).changedWorkspaces).toEqual([]) expect(await policyOf(sourceWorkspaceId)).toBe(true) + + // Each changed member's admins see the change in their own log, under that workspace's name. + const changed = new Map( + ( + await db + .select({ id: workspace.id, name: workspace.name }) + .from(workspace) + .where( + inArray( + workspace.id, + first.changedWorkspaces.map((member) => member.id) + ) + ) + ).map((member) => [member.id, member.name]) + ) + await vi.waitFor( + async () => { + const entries = await auditedFrom(childId) + // Exactly the changed members, once each: no missing, duplicate, or extra entry, + // including from the no-op repeat issued from the same workspace. + expect(entries.map((entry) => entry.resourceId).sort()).toEqual( + [...changed.keys()].sort() + ) + for (const entry of entries) { + expect(entry.workspaceId).toBe(entry.resourceId) + expect(entry.resourceName).toBe(changed.get(entry.resourceId!)) + expect(entry.metadata).toMatchObject({ + forkSyncNewWorkflowsExcluded: true, + originWorkspaceId: childId, + originWorkspaceName: changed.get(childId), + }) + } + }, + { timeout: 5000 } + ) expect(await excludedFor(sourceWorkflowId)).toBe(false) const [copy] = await db @@ -838,10 +892,26 @@ describe('authorized fork and sync against PostgreSQL', () => { ).toEqual([{ excluded: true }]) await db.update(workspace).set({ archivedAt: new Date() }).where(eq(workspace.id, childId)) - await setDefault(grandchildId, false) - expect(await policyOf(sourceWorkspaceId)).toBe(false) - expect(await policyOf(grandchildId)).toBe(false) + const fromGrandchild = await setDefault(grandchildId, false) + // Every live member flips back - the ones the first change covered and the two forks + // created since - while the archived child is neither written nor audited. + const expectedFromGrandchild = [ + ...[...changed.keys()].filter((id) => id !== childId), + newForkId, + grandchildId, + ].sort() + expect(fromGrandchild.changedWorkspaces.map((member) => member.id).sort()).toEqual( + expectedFromGrandchild + ) + for (const id of expectedFromGrandchild) expect(await policyOf(id)).toBe(false) expect(await policyOf(childId)).toBe(true) + await vi.waitFor( + async () => { + const entries = await auditedFrom(grandchildId) + expect(entries.map((entry) => entry.resourceId).sort()).toEqual(expectedFromGrandchild) + }, + { timeout: 5000 } + ) } finally { await db.update(workspace).set({ archivedAt: null }).where(eq(workspace.id, childId)) await setDefault(sourceWorkspaceId, false)