Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
59 changes: 58 additions & 1 deletion apps/sim/executor/execution/block-executor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
import { loggerMock } from '@sim/testing'
import { beforeEach, describe, expect, it, vi } from 'vitest'
import { clearLargeValueCacheForTests } from '@/lib/execution/payloads/cache'
import { createLargeArrayManifest } from '@/lib/execution/payloads/large-array-manifest'
import { isLargeArrayManifest } from '@/lib/execution/payloads/large-array-manifest-metadata'
import { isLargeValueRef } from '@/lib/execution/payloads/large-value-ref'
import { projectTraceSpansForSecrets } from '@/lib/logs/execution/trace-secret-projection'
Expand All @@ -24,8 +25,10 @@ const blockExecutorBaseLogger =
loggerMock.createLogger.mock.results[blockExecutorLoggerCallIndex]?.value
if (!blockExecutorBaseLogger) throw new Error('BlockExecutor logger mock was not initialized')

const { mockUploadFile } = vi.hoisted(() => ({
const { mockUploadFile, mockDownloadFile, mockMaskBatch } = vi.hoisted(() => ({
mockUploadFile: vi.fn(),
mockDownloadFile: vi.fn(),
mockMaskBatch: vi.fn(),
}))

vi.mock('@/ee/access-control/utils/permission-check', () => ({
Expand All @@ -35,9 +38,14 @@ vi.mock('@/ee/access-control/utils/permission-check', () => ({
vi.mock('@/lib/uploads', () => ({
StorageService: {
uploadFile: mockUploadFile,
downloadFile: mockDownloadFile,
},
}))

vi.mock('@/lib/guardrails/mask-client', () => ({
maskPIIBatchViaHttp: mockMaskBatch,
}))

vi.mock('@/lib/logs/execution/pii-redaction', async (importOriginal) => {
const actual = await importOriginal<typeof import('@/lib/logs/execution/pii-redaction')>()
return {
Expand Down Expand Up @@ -94,6 +102,55 @@ describe('BlockExecutor', () => {
mockUploadFile.mockImplementation(async ({ customKey }) => ({ key: customKey }))
})

it('redacts an authorized prior-execution manifest returned by a block under the current execution', async () => {
const items = [{ email: 'alice@example.com', count: 7 }]
const manifest = await createLargeArrayManifest(items, {
workspaceId: 'workspace-1',
workflowId: 'workflow-1',
executionId: 'source-execution',
})
clearLargeValueCacheForTests()
mockUploadFile.mockClear()
mockDownloadFile.mockResolvedValue(Buffer.from(JSON.stringify(items)))
mockMaskBatch.mockImplementation(async (texts: string[]) =>
texts.map((text) => text.replaceAll('alice@example.com', '<EMAIL_ADDRESS>'))
)
const block = createBlock()
const workflow: SerializedWorkflow = {
version: '1',
blocks: [block],
connections: [],
loops: {},
parallels: {},
}
const state = new ExecutionState()
const resolver = new VariableResolver(workflow, {}, state)
const handler: BlockHandler = {
canHandle: () => true,
execute: async () => ({ result: manifest }),
}
const executor = new BlockExecutor([handler], resolver, {}, state)
const ctx = createContext(state)
ctx.largeValueExecutionIds = ['source-execution']
ctx.piiBlockOutputRedaction = { enabled: true, entityTypes: ['EMAIL_ADDRESS'], language: 'en' }

await executor.execute(ctx, createNode(block), block)

expect(state.getBlockOutput(block.id)?.result).toMatchObject({
preview: [{ email: '<EMAIL_ADDRESS>', count: 7 }],
chunks: [{ ref: { executionId: 'execution-1' } }],
})
expect(mockDownloadFile).toHaveBeenCalledWith(
expect.objectContaining({ key: manifest.chunks[0].ref.key })
)
expect(mockUploadFile).toHaveBeenCalledWith(
expect.objectContaining({
customKey: expect.stringContaining('execution/workspace-1/workflow-1/execution-1/'),
file: Buffer.from(JSON.stringify([{ email: '<EMAIL_ADDRESS>', count: 7 }])),
})
)
})

it('persists function output arrays as manifests in execution state', async () => {
const block = createBlock()
const workflow: SerializedWorkflow = {
Expand Down
3 changes: 3 additions & 0 deletions apps/sim/executor/execution/block-executor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -358,6 +358,9 @@ export class BlockExecutor {
workspaceId: blockCtx.workspaceId,
workflowId: blockCtx.workflowId,
executionId: blockCtx.executionId,
largeValueExecutionIds: blockCtx.largeValueExecutionIds,
largeValueKeys: blockCtx.largeValueKeys,
allowLargeValueWorkflowScope: blockCtx.allowLargeValueWorkflowScope,
userId: blockCtx.userId,
},
})
Expand Down
196 changes: 194 additions & 2 deletions apps/sim/lib/workflows/executor/execution-core.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,14 @@ import {
workflowsUtilsMock,
workflowsUtilsMockFns,
} from '@sim/testing'
import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest'
import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import * as retention from '@/lib/billing/retention'
import { clearLargeValueCacheForTests } from '@/lib/execution/payloads/cache'
import type { LargeArrayManifest } from '@/lib/execution/payloads/large-array-manifest'
import type { LargeValueRef } from '@/lib/execution/payloads/large-value-ref'
import type { LoggingSession } from '@/lib/logs/execution/logging-session'
import { ExecutionSnapshot } from '@/executor/execution/snapshot'
import type { SerializableExecutionState } from '@/executor/execution/types'

const {
mergeSubblockStateWithValuesMock,
Expand All @@ -32,6 +39,9 @@ const {
projectDisplayContentMock,
projectDiagnosticErrorMock,
decryptSecretMock,
downloadFileMock,
uploadFileMock,
maskBatchMock,
} = vi.hoisted(() => ({
mergeSubblockStateWithValuesMock: vi.fn(),
safeStartMock: vi.fn(),
Expand All @@ -54,6 +64,9 @@ const {
projectDisplayContentMock: vi.fn(),
projectDiagnosticErrorMock: vi.fn(),
decryptSecretMock: vi.fn(),
downloadFileMock: vi.fn(),
uploadFileMock: vi.fn(),
maskBatchMock: vi.fn(),
}))

const getPersonalAndWorkspaceEnvMock = environmentUtilsMockFns.mockGetPersonalAndWorkspaceEnv
Expand All @@ -67,6 +80,19 @@ const loadWorkflowDeploymentVersionStateMock =
workflowsPersistenceUtilsMockFns.mockLoadWorkflowDeploymentVersionState
const updateWorkflowRunCountsMock = workflowsUtilsMockFns.mockUpdateWorkflowRunCounts

vi.mock('@/lib/uploads', () => ({
StorageService: { downloadFile: downloadFileMock, uploadFile: uploadFileMock },
}))

vi.mock('@/lib/guardrails/mask-client', () => ({
maskPIIBatchViaHttp: maskBatchMock,
}))

vi.mock('@/lib/execution/payloads/large-value-metadata', () => ({
registerLargeValueOwner: vi.fn().mockResolvedValue(true),
addLargeValueReference: vi.fn().mockResolvedValue(undefined),
}))

vi.mock('@/lib/execution/cancellation', () => ({
clearExecutionCancellation: clearExecutionCancellationMock,
}))
Expand Down Expand Up @@ -120,7 +146,7 @@ import {
executeWorkflowCore,
FINALIZED_EXECUTION_ID_TTL_MS,
wasExecutionFinalizedByCore,
} from './execution-core'
} from '@/lib/workflows/executor/execution-core'

const executionCoreLoggerCallIndex = loggerMock.createLogger.mock.calls.findIndex(
([name]) => name === 'ExecutionCore'
Expand Down Expand Up @@ -993,6 +1019,172 @@ describe('executeWorkflowCore terminal finalization sequencing', () => {
expect(loadWorkflowDeploymentVersionStateMock).not.toHaveBeenCalled()
})

describe('PII redaction of restored large values', () => {
const sourceItems = [{ email: 'alice@example.com', count: 7 }]
const maskedItems = [{ email: '<EMAIL_ADDRESS>', count: 7 }]
const sourceBytes = Buffer.from(JSON.stringify(sourceItems))

function createManifest(
workspaceId = 'workspace-1',
workflowId = 'workflow-1',
executionId = 'source-execution'
): LargeArrayManifest {
const ref: LargeValueRef = {
__simLargeValueRef: true,
version: 1,
id: 'lv_123456789012',
kind: 'array',
size: sourceBytes.length,
executionId,
key: `execution/${workspaceId}/${workflowId}/${executionId}/large-value-lv_123456789012.json`,
}
return {
__simLargeArrayManifest: true,
version: 2,
kind: 'array',
totalCount: 1,
chunkCount: 1,
byteSize: sourceBytes.length,
chunks: [{ ref, count: 1, byteSize: sourceBytes.length }],
preview: sourceItems,
}
}

function createRestoredState(manifest: LargeArrayManifest): SerializableExecutionState {
return {
blockStates: { previous: { output: { result: manifest } } },
executedBlocks: ['previous'],
blockLogs: [],
decisions: { router: {}, condition: {} },
completedLoops: [],
activeExecutionPath: [],
trustedLargeValueAccess: { executionIds: [], largeValueKeys: [], fileKeys: [] },
}
}

function createPiiSnapshot(state?: SerializableExecutionState, input: unknown = {}) {
const base = createSnapshot()
return new ExecutionSnapshot(
{ ...base.metadata, resumeFromSnapshot: state !== undefined },
base.workflow,
input,
{},
[],
state
)
}

beforeEach(() => {
clearLargeValueCacheForTests()
vi.spyOn(retention, 'resolveEffectivePiiRedaction').mockReturnValue({
...retention.DEFAULT_PII_REDACTION,
input: {
enabled: true,
entityTypes: ['EMAIL_ADDRESS'],
language: 'en',
customPatterns: [],
},
blockOutputs: {
enabled: true,
entityTypes: ['EMAIL_ADDRESS'],
language: 'en',
customPatterns: [],
},
})
downloadFileMock.mockResolvedValue(sourceBytes)
uploadFileMock.mockImplementation(async ({ customKey }: { customKey: string }) => ({
key: customKey,
}))
maskBatchMock.mockImplementation(async (texts: string[]) =>
texts.map((text) => text.replaceAll('alice@example.com', '<EMAIL_ADDRESS>'))
)
executorExecuteMock.mockResolvedValue({
success: true,
status: 'completed',
output: { done: true },
logs: [],
metadata: { duration: 1, startTime: 'start', endTime: 'end' },
})
})

afterEach(() => {
vi.restoreAllMocks()
clearLargeValueCacheForTests()
})

it.each(['source execution', 'trusted key', 'resume', 'input'] as const)(
'masks cached manifest content with %s access and stores it under the new execution',
async (mode) => {
const manifest = createManifest()
const state = createRestoredState(manifest)
if (mode === 'trusted key')
state.trustedLargeValueAccess!.largeValueKeys = [manifest.chunks[0].ref.key!]
const snapshot = createPiiSnapshot(
mode === 'resume' ? state : undefined,
mode === 'input' ? { result: manifest } : {}
)
if (mode === 'input') snapshot.metadata.largeValueExecutionIds = ['source-execution']
const result = await executeWorkflowCore({
snapshot,
callbacks: {},
loggingSession: loggingSession as unknown as LoggingSession,
...(mode === 'resume' || mode === 'input'
? {}
: {
runFromBlock: {
startBlockId: 'start-block',
sourceSnapshot: state,
sourceExecutionId:
mode === 'trusted key' ? 'intermediate-execution' : 'source-execution',
},
}),
})
await loggingSession.setPostExecutionPromise.mock.calls[0][0]
expect(result.success).toBe(true)
expect(executorExecuteMock).toHaveBeenCalledOnce()
expect(downloadFileMock).toHaveBeenCalledWith(
expect.objectContaining({ key: manifest.chunks[0].ref.key, maxBytes: 64 * 1024 * 1024 })
)
expect(uploadFileMock).toHaveBeenCalledWith(
expect.objectContaining({
customKey: expect.stringContaining('execution/workspace-1/workflow-1/execution-1/'),
file: Buffer.from(JSON.stringify(maskedItems)),
})
)
if (mode !== 'input')
expect(state.blockStates.previous.output).toMatchObject({
result: { preview: maskedItems },
})
}
)

it.each([
['another workspace', 'workspace-2', 'workflow-1', 'source-execution'],
['another workflow', 'workspace-1', 'workflow-2', 'source-execution'],
['an unauthorized execution', 'workspace-1', 'workflow-1', 'unrelated-execution'],
])(
'refuses cached manifest content from %s before reading storage',
async (_, workspaceId, workflowId, executionId) => {
const state = createRestoredState(createManifest(workspaceId, workflowId, executionId))
await expect(
executeWorkflowCore({
snapshot: createPiiSnapshot(),
callbacks: {},
loggingSession: loggingSession as unknown as LoggingSession,
runFromBlock: {
startBlockId: 'start-block',
sourceExecutionId: 'source-execution',
sourceSnapshot: state,
},
})
).rejects.toThrow('Large execution value is not available in this execution.')
expect(downloadFileMock).not.toHaveBeenCalled()
expect(uploadFileMock).not.toHaveBeenCalled()
expect(executorExecuteMock).not.toHaveBeenCalled()
}
)
})

it('marks inherited client run-from-block provenance incomplete', async () => {
executorExecuteMock.mockResolvedValue({
success: true,
Expand Down
6 changes: 6 additions & 0 deletions apps/sim/lib/workflows/executor/execution-core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -888,6 +888,9 @@ async function executeWorkflowCoreImpl(
workspaceId: providedWorkspaceId,
workflowId,
executionId,
largeValueExecutionIds,
largeValueKeys,
allowLargeValueWorkflowScope,
userId: userId ?? undefined,
},
})
Expand Down Expand Up @@ -918,6 +921,9 @@ async function executeWorkflowCoreImpl(
workspaceId: providedWorkspaceId,
workflowId,
executionId,
largeValueExecutionIds,
largeValueKeys,
allowLargeValueWorkflowScope,
userId: userId ?? undefined,
},
}
Expand Down
Loading