diff --git a/apps/sim/app/api/table/row-secret-provenance.test.ts b/apps/sim/app/api/table/row-secret-provenance.test.ts index 21cbf24d0a5..6c4f66d4d56 100644 --- a/apps/sim/app/api/table/row-secret-provenance.test.ts +++ b/apps/sim/app/api/table/row-secret-provenance.test.ts @@ -3,7 +3,6 @@ */ import { createMockRequest } from '@sim/testing' import { describe, expect, it } from 'vitest' -import { AuthType } from '@/lib/auth/hybrid' import { PRIVATE_SECRET_PROVENANCE_BUNDLE_V1, PRIVATE_SECRET_PROVENANCE_FIELD, @@ -13,29 +12,16 @@ import { RESOLVED_SECRET_PROVENANCE_METADATA_V1, } from '@/lib/execution/private-tool-metadata' import { TableRowProvenanceError } from '@/lib/table/application/row-secret-provenance' -import { rowDataNameToId } from '@/lib/table/column-keys' import { tableRowSecretProvenanceSelectionKey } from '@/lib/table/secret-provenance-selection' -import type { RowData } from '@/lib/table/types' import { - createTableWriteProvenanceTargets, finalizeTableRowsProvenance, negotiateTableRowsProvenance, readTableRowProvenanceEnvelope, - resolveTableWriteSecretProvenance, } from '@/app/api/table/row-secret-provenance' const USER_ID = 'user-1' const WORKSPACE_ID = 'ws-1' -/** Mirrors the internal-JWT wire translator: names → ids, unknown names dropped. */ -const ID_BY_NAME = new Map([ - ['email', 'col_email'], - ['company', 'col_company'], -]) - -const translateNames = (data: RowData): RowData => rowDataNameToId(data, ID_BY_NAME) -const translateIdentity = (data: RowData): RowData => data - function traceProvenance() { return { version: 1, @@ -59,221 +45,6 @@ function bundleRequest(selectionKeys: string[]) { return { request, payload } } -describe('createTableWriteProvenanceTargets', () => { - it('maps column names to their storage ids', () => { - const targets = createTableWriteProvenanceTargets([{ email: 'a@b.c' }], translateNames) - - expect(targets).toEqual([ - { - selectionKey: tableRowSecretProvenanceSelectionKey(0, 'email'), - rowKey: '0', - columnId: 'col_email', - }, - ]) - }) - - it('returns a null column id for a column the wire translator drops', () => { - const targets = createTableWriteProvenanceTargets( - [{ email: 'a@b.c', notAColumn: 'x' }], - translateNames - ) - - expect(targets).toHaveLength(2) - expect(targets[0].columnId).toBe('col_email') - expect(targets[1]).toEqual({ - selectionKey: tableRowSecretProvenanceSelectionKey(0, 'notAColumn'), - rowKey: '0', - columnId: null, - }) - }) - - it('keeps one target per submitted column so bundle selections stay paired', () => { - const targets = createTableWriteProvenanceTargets( - [{ notAColumn: 'x', alsoNotAColumn: 'y' }], - translateNames - ) - - expect(targets.map((target) => target.columnId)).toEqual([null, null]) - }) - - it('passes column ids through for identity (session) translation', () => { - const targets = createTableWriteProvenanceTargets([{ col_email: 'a@b.c' }], translateIdentity) - - expect(targets[0].columnId).toBe('col_email') - }) - - it('keys targets by row index across multiple rows', () => { - const targets = createTableWriteProvenanceTargets( - [{ email: 'a@b.c' }, { company: 'Acme' }], - translateNames - ) - - expect(targets.map((target) => target.rowKey)).toEqual(['0', '1']) - expect(targets[1].selectionKey).toBe(tableRowSecretProvenanceSelectionKey(1, 'company')) - }) -}) - -describe('resolveTableWriteSecretProvenance', () => { - it('records no provenance for a dropped column on an unsupported session write', () => { - const rows = [{ email: 'a@b.c', notAColumn: 'x' }] - const result = resolveTableWriteSecretProvenance({ - request: createMockRequest('POST', { rows }), - payload: { rows }, - authType: AuthType.SESSION, - userId: USER_ID, - workspaceId: WORKSPACE_ID, - targets: createTableWriteProvenanceTargets(rows, translateNames), - rowKeys: ['0'], - }) - - expect(result.success).toBe(true) - if (!result.success) return - expect(Object.keys(result.provenanceByRowKey?.['0'].columns ?? {})).toEqual(['col_email']) - }) - - it('accepts a complete bundle that covers a dropped column', () => { - const rows = [{ email: 'a@b.c', notAColumn: 'x' }] - const { request, payload } = bundleRequest([ - tableRowSecretProvenanceSelectionKey(0, 'email'), - tableRowSecretProvenanceSelectionKey(0, 'notAColumn'), - ]) - - const result = resolveTableWriteSecretProvenance({ - request, - payload, - authType: AuthType.INTERNAL_JWT, - userId: USER_ID, - workspaceId: WORKSPACE_ID, - targets: createTableWriteProvenanceTargets(rows, translateNames), - rowKeys: ['0'], - }) - - expect(result.success).toBe(true) - if (!result.success) return - expect(Object.keys(result.provenanceByRowKey?.['0'].columns ?? {})).toEqual(['col_email']) - }) - - it('stores provenance for a fully translatable bundle', () => { - const rows = [{ email: 'a@b.c', company: 'Acme' }] - const { request, payload } = bundleRequest([ - tableRowSecretProvenanceSelectionKey(0, 'email'), - tableRowSecretProvenanceSelectionKey(0, 'company'), - ]) - - const result = resolveTableWriteSecretProvenance({ - request, - payload, - authType: AuthType.INTERNAL_JWT, - userId: USER_ID, - workspaceId: WORKSPACE_ID, - targets: createTableWriteProvenanceTargets(rows, translateNames), - rowKeys: ['0'], - }) - - expect(result.success).toBe(true) - if (!result.success) return - expect(Object.keys(result.provenanceByRowKey?.['0'].columns ?? {}).sort()).toEqual([ - 'col_company', - 'col_email', - ]) - }) - - it('rejects a bundle whose selection matches no submitted column', () => { - const rows = [{ email: 'a@b.c' }] - const { request, payload } = bundleRequest([tableRowSecretProvenanceSelectionKey(0, 'company')]) - - const result = resolveTableWriteSecretProvenance({ - request, - payload, - authType: AuthType.INTERNAL_JWT, - userId: USER_ID, - workspaceId: WORKSPACE_ID, - targets: createTableWriteProvenanceTargets(rows, translateNames), - rowKeys: ['0'], - }) - - expect(result.success).toBe(false) - }) - - it('accepts a different source user in the authorized destination workspace', () => { - const rows = [{ email: 'a@b.c' }] - const payload = { - [PRIVATE_SECRET_PROVENANCE_FIELD]: { - version: 1, - complete: true, - selections: [ - { - key: tableRowSecretProvenanceSelectionKey(0, 'email'), - provenance: { - ...traceProvenance(), - scope: { userId: 'someone-else', workspaceId: WORKSPACE_ID }, - }, - }, - ], - }, - } - - const result = resolveTableWriteSecretProvenance({ - request: createMockRequest('POST', payload, { - [PRIVATE_SECRET_PROVENANCE_HEADER]: PRIVATE_SECRET_PROVENANCE_BUNDLE_V1, - }), - payload, - authType: AuthType.INTERNAL_JWT, - userId: USER_ID, - workspaceId: WORKSPACE_ID, - targets: createTableWriteProvenanceTargets(rows, translateNames), - rowKeys: ['0'], - }) - - expect(result.success).toBe(true) - if (!result.success) return - expect(result.provenanceByRowKey?.['0'].columns.col_email).toMatchObject({ - scope: { userId: 'someone-else', workspaceId: WORKSPACE_ID }, - }) - }) - - it('rejects a bundle whose selection comes from another workspace', () => { - const rows = [{ email: 'a@b.c' }] - const payload = { - [PRIVATE_SECRET_PROVENANCE_FIELD]: { - version: 1, - complete: true, - selections: [ - { - key: tableRowSecretProvenanceSelectionKey(0, 'email'), - provenance: { - ...traceProvenance(), - scope: { userId: USER_ID, workspaceId: 'another-workspace' }, - }, - }, - ], - }, - } - - const result = resolveTableWriteSecretProvenance({ - request: createMockRequest('POST', payload, { - [PRIVATE_SECRET_PROVENANCE_HEADER]: PRIVATE_SECRET_PROVENANCE_BUNDLE_V1, - }), - payload, - authType: AuthType.INTERNAL_JWT, - userId: USER_ID, - workspaceId: WORKSPACE_ID, - targets: createTableWriteProvenanceTargets(rows, translateNames), - rowKeys: ['0'], - }) - - expect(result.success).toBe(false) - }) -}) - -/** - * The transport half of the envelope, used by the migrated single-row routes. - * - * These are the only cover these helpers have: the route tests assert `mapInput` - * and `present`, so each helper could be replaced by a constant without a route - * test noticing — and a constant `readTableRowProvenanceEnvelope` would silently - * downgrade every executor write from a stamped bundle to untracked. - */ describe('readTableRowProvenanceEnvelope', () => { it('reports no envelope when the caller sent none', () => { const request = createMockRequest('PATCH', { data: {} }) diff --git a/apps/sim/app/api/table/row-secret-provenance.ts b/apps/sim/app/api/table/row-secret-provenance.ts index d455a6617d8..583fc4a66fb 100644 --- a/apps/sim/app/api/table/row-secret-provenance.ts +++ b/apps/sim/app/api/table/row-secret-provenance.ts @@ -1,10 +1,5 @@ -import { type NextRequest, NextResponse } from 'next/server' -import { AuthType, type AuthTypeValue } from '@/lib/auth/hybrid' -import { isPrivateSecretProvenanceScopeCompatible } from '@/lib/execution/durable-secret-provenance' -import { - inspectPrivateSecretProvenanceRequest, - isPrivateSecretProvenanceBundleV1, -} from '@/lib/execution/model-input-provenance' +import type { NextRequest } from 'next/server' +import { inspectPrivateSecretProvenanceRequest } from '@/lib/execution/model-input-provenance' import { negotiatePrivateToolMetadataResponse, RESOLVED_SECRET_PROVENANCE_METADATA_V1, @@ -14,191 +9,6 @@ import { type TableRowProvenanceEnvelope, TableRowProvenanceError, } from '@/lib/table/application/row-secret-provenance' -import { loadTableRowSecretProvenance } from '@/lib/table/rows/secret-provenance' -import { tableRowSecretProvenanceSelectionKey } from '@/lib/table/secret-provenance-selection' -import type { RowData, TableRowSecretProvenanceWrite } from '@/lib/table/types' - -interface TableRowCrossing { - id: string - updatedAt: Date | string - data?: RowData -} - -type TableWriteProvenanceResult = - | { - success: true - provenanceByRowKey: Record | undefined - } - | { success: false; response: NextResponse } - -interface TableWriteProvenanceTarget { - selectionKey: string - rowKey: string - /** Storage column id, or `null` when the wire translator drops this column. */ - columnId: string | null -} - -/** - * Maps tool-facing column names to the stable storage ids used by the sidecar. - * - * The wire translator drops keys that name no column in the table schema, and the - * write path drops them identically, so such a column is simply never persisted. - * It still gets a target — callers key one provenance selection per column they - * sent, and the completeness check pairs the two — but with a `null` column id so - * no provenance is recorded for a value that was never stored. - */ -export function createTableWriteProvenanceTargets( - rows: readonly RowData[], - translate: (data: RowData) => RowData -): TableWriteProvenanceTarget[] { - return rows.flatMap((row, rowIndex) => - Object.entries(row).map(([columnKey, value]) => { - const translatedKeys = Object.keys(translate({ [columnKey]: value })) - return { - selectionKey: tableRowSecretProvenanceSelectionKey(rowIndex, columnKey), - rowKey: String(rowIndex), - columnId: translatedKeys.length === 1 ? translatedKeys[0] : null, - } - }) - ) -} - -function invalidProvenanceResponse(): NextResponse { - return NextResponse.json({ error: 'Invalid table row secret provenance' }, { status: 400 }) -} - -/** - * Authenticates the private provenance envelope on a table mutation. Interactive - * writes receive an exact-empty report, internal legacy callers remain untracked, - * and only an internal JWT may submit encrypted provenance. - */ -export function resolveTableWriteSecretProvenance(options: { - request: NextRequest - payload: unknown - authType: AuthTypeValue | undefined - userId: string - workspaceId: string - targets: TableWriteProvenanceTarget[] - rowKeys: string[] -}): TableWriteProvenanceResult { - const { request } = options - const inspection = inspectPrivateSecretProvenanceRequest(request.headers, options.payload) - if (inspection.status === 'unsupported') { - if (options.authType === AuthType.INTERNAL_JWT) { - return { success: true, provenanceByRowKey: undefined } - } - const provenanceByRowKey: Record = {} - for (const rowKey of options.rowKeys) { - provenanceByRowKey[rowKey] = { complete: true, columns: {} } - } - for (const target of options.targets) { - if (target.columnId === null) continue - const row = provenanceByRowKey[target.rowKey] ?? { complete: true, columns: {} } - row.columns[target.columnId] = { - version: 1, - complete: true, - entries: [], - scope: { userId: options.userId, workspaceId: options.workspaceId }, - } - provenanceByRowKey[target.rowKey] = row - } - return { - success: true, - provenanceByRowKey, - } - } - if ( - inspection.status !== 'verified' || - options.authType !== AuthType.INTERNAL_JWT || - !isPrivateSecretProvenanceBundleV1(inspection.value) - ) { - return { success: false, response: invalidProvenanceResponse() } - } - - const bundle = inspection.value - const targetBySelectionKey = new Map( - options.targets.map((target) => [target.selectionKey, target]) - ) - if ( - bundle.complete && - (bundle.selections.length !== options.targets.length || - bundle.selections.some((selection) => !targetBySelectionKey.has(selection.key))) - ) { - return { success: false, response: invalidProvenanceResponse() } - } - const rowKeys = [...new Set(options.rowKeys)] - if (!bundle.complete) { - return { - success: true, - provenanceByRowKey: Object.fromEntries( - rowKeys.map((rowKey) => [rowKey, { complete: false, columns: {} }]) - ), - } - } - - const provenanceByRowKey: Record = Object.fromEntries( - rowKeys.map((rowKey) => [rowKey, { complete: true, columns: {} }]) - ) - for (const selection of bundle.selections) { - const target = targetBySelectionKey.get(selection.key) - if ( - !target || - !isPrivateSecretProvenanceScopeCompatible(selection.provenance.scope, { - userId: options.userId, - workspaceId: options.workspaceId, - }) - ) { - return { success: false, response: invalidProvenanceResponse() } - } - if (target.columnId === null) continue - if (Object.hasOwn(provenanceByRowKey[target.rowKey].columns, target.columnId)) { - return { success: false, response: invalidProvenanceResponse() } - } - provenanceByRowKey[target.rowKey].columns[target.columnId] = selection.provenance - } - return { success: true, provenanceByRowKey } -} - -/** - * Adds row provenance only for an authenticated internal caller that explicitly - * requested the supported private capability. All ordinary API/UI responses keep - * their existing wire shape. - */ -export async function createTableRowsResponse(options: { - request: NextRequest - authType: AuthTypeValue | undefined - userId: string - workspaceId: string - body: Record - rows: TableRowCrossing[] -}): Promise { - const { request } = options - const negotiation = negotiatePrivateToolMetadataResponse( - request.headers, - RESOLVED_SECRET_PROVENANCE_METADATA_V1, - options.authType === AuthType.INTERNAL_JWT - ) - if (negotiation.status === 'not-requested') return NextResponse.json(options.body) - if (negotiation.status === 'rejected') return invalidProvenanceResponse() - - const provenance = await loadTableRowSecretProvenance( - options.rows.map((row) => ({ - id: row.id, - updatedAt: row.updatedAt, - ...(row.data ? { selectedValues: row.data } : {}), - })), - { - userId: options.userId, - workspaceId: options.workspaceId, - } - ) - const envelope = serializePrivateToolMetadataResponseEnvelope( - options.body, - RESOLVED_SECRET_PROVENANCE_METADATA_V1, - provenance - ) - return NextResponse.json(envelope.body, { headers: envelope.headers }) -} /** * Reads the private provenance envelope off the request. Transport only — the diff --git a/apps/sim/background/dispatch-cancel-guard.test.ts b/apps/sim/background/dispatch-cancel-guard.test.ts index ecf287abb5c..6bf8b9cb24f 100644 --- a/apps/sim/background/dispatch-cancel-guard.test.ts +++ b/apps/sim/background/dispatch-cancel-guard.test.ts @@ -25,6 +25,7 @@ vi.mock('@/lib/table/dispatcher', () => ({ vi.mock('@/lib/table/service', () => ({ getTableById: mocks.getTableById })) vi.mock('@/lib/table/rows/service', () => ({ getRowById: mocks.getRowById, + getRowSummaryById: vi.fn(), updateRow: vi.fn(), })) vi.mock('@/lib/workflows/executor/execute-workflow', () => ({ diff --git a/apps/sim/background/drain-governed-subject.test.ts b/apps/sim/background/drain-governed-subject.test.ts index cb3d28654ca..6e826a1e668 100644 --- a/apps/sim/background/drain-governed-subject.test.ts +++ b/apps/sim/background/drain-governed-subject.test.ts @@ -7,6 +7,8 @@ import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' const mocks = vi.hoisted(() => ({ getTableById: vi.fn(), getRowById: vi.fn(), + getRowSummaryById: vi.fn(), + createProvenanceReader: vi.fn(), pickNextEligibleGroupForRow: vi.fn(), writeWorkflowGroupState: vi.fn(), markWorkflowGroupPickedUp: vi.fn(), @@ -14,11 +16,15 @@ const mocks = vi.hoisted(() => ({ getEnrichment: vi.fn(), readStampedCapabilitySubject: vi.fn(), checkAttributedUsageLimits: vi.fn(), - loadTableRowSecretProvenance: vi.fn(), + exportProvenance: vi.fn(), })) vi.mock('@/lib/table/service', () => ({ getTableById: mocks.getTableById })) -vi.mock('@/lib/table/rows/service', () => ({ getRowById: mocks.getRowById, updateRow: vi.fn() })) +vi.mock('@/lib/table/rows/service', () => ({ + getRowById: mocks.getRowById, + getRowSummaryById: mocks.getRowSummaryById, + updateRow: vi.fn(), +})) vi.mock('@/lib/table/rows/executions', () => ({ readStampedCapabilitySubject: mocks.readStampedCapabilitySubject, })) @@ -49,7 +55,12 @@ vi.mock('@/lib/billing/core/billing-attribution', () => ({ vi.mock('@/lib/table/rows/secret-provenance', () => ({ createExactEmptyTableRowSecretProvenance: () => ({ complete: true, columns: {} }), createTableRowSecretProvenanceFromRegistry: () => ({ complete: true, columns: {} }), - loadTableRowSecretProvenance: mocks.loadTableRowSecretProvenance, + TableRowProvenanceReader: class { + constructor(scope: unknown, selectedColumnIds: unknown) { + mocks.createProvenanceReader(scope, selectedColumnIds) + } + exportProvenance = mocks.exportProvenance + }, })) vi.mock('@/lib/table/events', () => ({ appendTableEvent: vi.fn() })) @@ -142,6 +153,9 @@ describe('draining another dispatch’s pre-stamped marker', () => { beforeEach(() => { vi.clearAllMocks() resetDbChainMock() + mocks.getRowSummaryById.mockImplementation((tableId, rowId, workspaceId) => + mocks.getRowById(tableId, rowId, workspaceId) + ) mocks.getTableById.mockResolvedValue(TABLE) mocks.getEnrichment.mockReturnValue({ id: 'enrich-1', @@ -152,7 +166,7 @@ describe('draining another dispatch’s pre-stamped marker', () => { mocks.checkAttributedUsageLimits.mockResolvedValue({ isExceeded: false }) mocks.markWorkflowGroupPickedUp.mockResolvedValue('picked-up') mocks.writeWorkflowGroupState.mockResolvedValue('wrote') - mocks.loadTableRowSecretProvenance.mockResolvedValue({ + mocks.exportProvenance.mockReturnValue({ scope: { userId: 'user-1', workspaceId: 'workspace-1' }, byRowId: {}, }) diff --git a/apps/sim/background/enrichment-capability-subject.test.ts b/apps/sim/background/enrichment-capability-subject.test.ts index a89951e86fd..d992ed749f9 100644 --- a/apps/sim/background/enrichment-capability-subject.test.ts +++ b/apps/sim/background/enrichment-capability-subject.test.ts @@ -11,6 +11,8 @@ const mocks = vi.hoisted(() => ({ })), getTableById: vi.fn(), getRowById: vi.fn(), + getRowSummaryById: vi.fn(), + createProvenanceReader: vi.fn(), updateRow: vi.fn(), pickNextEligibleGroupForRow: vi.fn(), stashCellContextForResume: vi.fn(), @@ -23,7 +25,7 @@ const mocks = vi.hoisted(() => ({ runEnrichment: vi.fn(), skippedEnrichmentDetail: vi.fn(() => ({})), checkAttributedUsageLimits: vi.fn(async () => ({ isExceeded: false })), - loadTableRowSecretProvenance: vi.fn(async () => ({ scope: null, entries: [] })), + exportProvenance: vi.fn(() => ({ scope: null, entries: [] })), })) vi.mock('@/lib/core/network/config.server', () => ({ @@ -37,6 +39,7 @@ vi.mock('@/lib/workspaces/application/workspace-context', () => ({ vi.mock('@/lib/table/service', () => ({ getTableById: mocks.getTableById })) vi.mock('@/lib/table/rows/service', () => ({ getRowById: mocks.getRowById, + getRowSummaryById: mocks.getRowSummaryById, updateRow: mocks.updateRow, })) vi.mock('@/lib/table/cell-write', () => ({ @@ -61,7 +64,12 @@ vi.mock('@/lib/billing/core/billing-attribution', () => ({ vi.mock('@/lib/table/rows/secret-provenance', () => ({ createExactEmptyTableRowSecretProvenance: vi.fn(() => undefined), createTableRowSecretProvenanceFromRegistry: vi.fn(() => undefined), - loadTableRowSecretProvenance: mocks.loadTableRowSecretProvenance, + TableRowProvenanceReader: class { + constructor(scope: unknown, selectedColumnIds: unknown) { + mocks.createProvenanceReader(scope, selectedColumnIds) + } + exportProvenance = mocks.exportProvenance + }, })) vi.mock('@/executor/utils/resolved-secret-trace-registry', () => ({ ResolvedSecretTraceRegistry: class { @@ -156,6 +164,9 @@ describe('enrichment cell capability subject', () => { beforeEach(() => { vi.clearAllMocks() resetDbChainMock() + mocks.getRowSummaryById.mockImplementation((tableId, rowId, workspaceId) => + mocks.getRowById(tableId, rowId, workspaceId) + ) mocks.loadWorkspaceApplicationContext.mockResolvedValue({ workspaceOrganizationId: null }) mocks.getTableById.mockResolvedValue(TABLE) mocks.getRowById.mockResolvedValue({ @@ -176,6 +187,49 @@ describe('enrichment cell capability subject', () => { mocks.runEnrichment.mockResolvedValue({ result: {}, cost: 0, detail: {} }) }) + it('builds enrichment inputs from the captured row after pickup and excludes own outputs', async () => { + mocks.getTableById.mockResolvedValue({ + ...TABLE, + schema: { + ...TABLE.schema, + workflowGroups: [ + { + ...GROUP, + inputMappings: [ + ...GROUP.inputMappings, + { columnName: 'col-out', inputName: 'excluded' }, + ], + }, + ], + }, + }) + mocks.getRowSummaryById.mockResolvedValue({ + id: 'row-1', + data: { 'col-in': 'fresh.example.com', 'col-out': 'secret-output' }, + updatedAt: new Date('2026-09-14T00:00:00Z'), + }) + + await runRowCascadeLoop(payload('acting-user', 'acting-user') as never) + + expect(mocks.getRowSummaryById).toHaveBeenCalledWith( + 'table-1', + 'row-1', + 'workspace-1', + expect.any(Object) + ) + expect(mocks.createProvenanceReader).toHaveBeenCalledExactlyOnceWith( + { userId: 'acting-user', workspaceId: 'workspace-1' }, + new Set(['col-in']) + ) + expect(mocks.runEnrichment.mock.calls[0][1]).toEqual({ domain: 'fresh.example.com' }) + expect(mocks.markWorkflowGroupPickedUp.mock.invocationCallOrder[0]).toBeLessThan( + mocks.getRowSummaryById.mock.invocationCallOrder[0] + ) + expect(mocks.getRowSummaryById.mock.invocationCallOrder[0]).toBeLessThan( + mocks.runEnrichment.mock.invocationCallOrder[0] + ) + }) + it.each(['org_reserved', 'org_other', null])( 'restores the current workspace owner %s before running a queued enrichment', async (organizationId) => { diff --git a/apps/sim/background/workflow-column-execution.ts b/apps/sim/background/workflow-column-execution.ts index ef8c5022cf7..cf3e7ac0c71 100644 --- a/apps/sim/background/workflow-column-execution.ts +++ b/apps/sim/background/workflow-column-execution.ts @@ -40,7 +40,7 @@ import { appendTableEvent } from '@/lib/table/events' import { createExactEmptyTableRowSecretProvenance, createTableRowSecretProvenanceFromRegistry, - loadTableRowSecretProvenance, + TableRowProvenanceReader, } from '@/lib/table/rows/secret-provenance' import type { RowData, @@ -465,7 +465,7 @@ async function runWorkflowAndWriteTerminal( try { return await runWithRequestContext({ requestId }, async () => { - const { getRowById } = await import('@/lib/table/rows/service') + const { getRowById, getRowSummaryById } = await import('@/lib/table/rows/service') const { executeWorkflow } = await import('@/lib/workflows/executor/execute-workflow') const { loadDeployedWorkflowState } = await import('@/lib/workflows/persistence/utils') const { @@ -613,44 +613,57 @@ async function runWorkflowAndWriteTerminal( }) if (pickedUp === 'skipped') return 'error' - // Map table columns → enrichment input ids (skip this group's own outputs). - // `columnName` holds a column id; the mapper resolves select ids to names. - const ownOutputColumns = new Set(group.outputs.map((o) => o.columnName)) - const enrichmentInputMappings = (group.inputMappings ?? []).filter( - (mapping) => !ownOutputColumns.has(mapping.columnName) - ) - const enrichInputs = mapInputValues(row.data, table.schema.columns, enrichmentInputMappings) - - // Skip (don't error) rows missing a required input — common when a table - // is partially filled. Clear any prior output values so a stale result - // doesn't linger (and doesn't mark the group `completed`-and-filled, which - // would block the auto cascade from re-enriching once inputs return). - const isEmpty = isEmptyCellValue - const missingRequired = enrichment.inputs.some( - (i) => i.required && isEmpty(enrichInputs[i.id]) - ) - if (missingRequired) { - const clearPatch: RowData = {} - for (const out of group.outputs) { - if (!isEmpty(row.data[out.columnName])) clearPatch[out.columnName] = '' + try { + // Map table columns → enrichment input ids (skip this group's own outputs). + // `columnName` holds a column id; the mapper resolves select ids to names. + const ownOutputColumns = new Set(group.outputs.map((o) => o.columnName)) + const enrichmentInputMappings = (group.inputMappings ?? []).filter( + (mapping) => !ownOutputColumns.has(mapping.columnName) + ) + const readProvenance = new TableRowProvenanceReader( + { userId: enrichmentBillingAttribution.actorUserId, workspaceId }, + new Set(enrichmentInputMappings.map((mapping) => mapping.columnName)) + ) + const inputSource = await getRowSummaryById(tableId, rowId, workspaceId, readProvenance) + if (!inputSource) { + logger.warn(`Row ${rowId} vanished before enrichment input could be read`) + return 'error' } - await writeState( - { - status: 'completed', - executionId, - jobId: null, - workflowId: statusId, - error: null, - enrichmentDetails: skippedEnrichmentDetail(enrichment), - }, - clearPatch, - undefined, - createExactEmptyTableRowSecretProvenance(clearPatch) + const enrichInputs = mapInputValues( + inputSource.data, + table.schema.columns, + enrichmentInputMappings ) - return 'completed' - } - try { + // Skip (don't error) rows missing a required input — common when a table + // is partially filled. Clear any prior output values so a stale result + // doesn't linger (and doesn't mark the group `completed`-and-filled, which + // would block the auto cascade from re-enriching once inputs return). + const isEmpty = isEmptyCellValue + const missingRequired = enrichment.inputs.some( + (i) => i.required && isEmpty(enrichInputs[i.id]) + ) + if (missingRequired) { + const clearPatch: RowData = {} + for (const out of group.outputs) { + if (!isEmpty(inputSource.data[out.columnName])) clearPatch[out.columnName] = '' + } + await writeState( + { + status: 'completed', + executionId, + jobId: null, + workflowId: statusId, + error: null, + enrichmentDetails: skippedEnrichmentDetail(enrichment), + }, + clearPatch, + undefined, + createExactEmptyTableRowSecretProvenance(clearPatch) + ) + return 'completed' + } + if (attemptSignal.aborted) { await writeState({ ...buildTableAbortState({ @@ -663,21 +676,7 @@ async function runWorkflowAndWriteTerminal( }) return 'error' } - const inputProvenance = await loadTableRowSecretProvenance( - [ - { - id: row.id, - updatedAt: row.updatedAt, - selectedValues: Object.fromEntries( - enrichmentInputMappings.map((mapping) => [ - mapping.columnName, - row.data[mapping.columnName], - ]) - ), - }, - ], - { userId: enrichmentBillingAttribution.actorUserId, workspaceId } - ) + const inputProvenance = readProvenance.exportProvenance() const enrichmentRegistry = new ResolvedSecretTraceRegistry([], inputProvenance.scope) await enrichmentRegistry.importCrossingProvenance(inputProvenance, enrichInputs, { trusted: true, @@ -1015,10 +1014,26 @@ async function runWorkflowAndWriteTerminal( const inputColumns = table.schema.columns.filter( (c) => !ownOutputColumnIds.has(getColumnId(c)) ) + const inputMappings = group.inputMappings ?? [] + const readProvenance = new TableRowProvenanceReader( + { userId: workflowRecord.userId, workspaceId }, + new Set([ + ...inputColumns.map(getColumnId), + ...inputMappings.map((mapping) => mapping.columnName), + ]) + ) + const inputSource = await getRowSummaryById(tableId, rowId, workspaceId, readProvenance) + if (!inputSource) { + logger.warn(`Row ${rowId} vanished before workflow input could be read`) + return 'error' + } // One column list drives both the row and its headers so they cannot drift. // The mapper also resolves select option ids to names — the workflow author // sees "Open", not `opt_a1b2`. - const inputRow = fillMissingColumns(namedRowMapper(inputColumns)(row.data), inputColumns) + const inputRow = fillMissingColumns( + namedRowMapper(inputColumns)(inputSource.data), + inputColumns + ) const headers = inputColumns.map((c) => c.name) // When the group has explicit input mappings, feed the workflow's @@ -1026,8 +1041,7 @@ async function runWorkflowAndWriteTerminal( // Otherwise fall back to spreading every non-output column by name, so a // Start field still resolves when it matches a column name. `row`/`rawRow` // always carry the full (name-keyed) row for downstream reference. - const inputMappings = group.inputMappings ?? [] - const mappedInputs = mapInputValues(row.data, table.schema.columns, inputMappings) + const mappedInputs = mapInputValues(inputSource.data, table.schema.columns, inputMappings) const input = { ...(inputMappings.length > 0 ? mappedInputs : inputRow), @@ -1041,21 +1055,7 @@ async function runWorkflowAndWriteTerminal( tableName, timestamp: new Date().toISOString(), } - const rowInputProvenance = await loadTableRowSecretProvenance( - [ - { - id: row.id, - updatedAt: row.updatedAt, - selectedValues: Object.fromEntries( - inputColumns.map((column) => { - const columnId = getColumnId(column) - return [columnId, row.data[columnId]] - }) - ), - }, - ], - { userId: workflowRecord.userId, workspaceId } - ) + const rowInputProvenance = readProvenance.exportProvenance() const inputRegistry = new ResolvedSecretTraceRegistry([], rowInputProvenance.scope) await inputRegistry.importCrossingProvenance(rowInputProvenance, input, { trusted: true }) diff --git a/apps/sim/background/workflow-group-governed-subject.test.ts b/apps/sim/background/workflow-group-governed-subject.test.ts index 5c53a5ad304..b9e72c77717 100644 --- a/apps/sim/background/workflow-group-governed-subject.test.ts +++ b/apps/sim/background/workflow-group-governed-subject.test.ts @@ -7,6 +7,8 @@ import { beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' const mocks = vi.hoisted(() => ({ getTableById: vi.fn(), getRowById: vi.fn(), + getRowSummaryById: vi.fn(), + createProvenanceReader: vi.fn(), pickNextEligibleGroupForRow: vi.fn(), stashCellContextForResume: vi.fn(), writeWorkflowGroupState: vi.fn(), @@ -15,13 +17,16 @@ const mocks = vi.hoisted(() => ({ loadDeployedWorkflowState: vi.fn(), executeWorkflow: vi.fn(), preprocessExecution: vi.fn(), - loadTableRowSecretProvenance: vi.fn(), + exportProvenance: vi.fn(), findStartBlock: vi.fn(), + flattenWorkflowOutputs: vi.fn(), + normalizeInputFormatValue: vi.fn(), })) vi.mock('@/lib/table/service', () => ({ getTableById: mocks.getTableById })) vi.mock('@/lib/table/rows/service', () => ({ getRowById: mocks.getRowById, + getRowSummaryById: mocks.getRowSummaryById, updateRow: vi.fn(), })) vi.mock('@/lib/table/workflow-columns', () => ({ @@ -58,8 +63,12 @@ vi.mock('@/lib/workflows/executor/execute-workflow', () => ({ vi.mock('@/lib/workflows/triggers/triggers', () => ({ TriggerUtils: { findStartBlock: mocks.findStartBlock }, })) -vi.mock('@/lib/workflows/blocks/flatten-outputs', () => ({ flattenWorkflowOutputs: () => [] })) -vi.mock('@/lib/workflows/input-format', () => ({ normalizeInputFormatValue: () => [] })) +vi.mock('@/lib/workflows/blocks/flatten-outputs', () => ({ + flattenWorkflowOutputs: mocks.flattenWorkflowOutputs, +})) +vi.mock('@/lib/workflows/input-format', () => ({ + normalizeInputFormatValue: mocks.normalizeInputFormatValue, +})) vi.mock('@/lib/execution/preprocessing', () => ({ preprocessExecution: mocks.preprocessExecution, })) @@ -69,7 +78,12 @@ vi.mock('@/lib/table/admission-retry', () => ({ vi.mock('@/lib/table/rows/secret-provenance', () => ({ createExactEmptyTableRowSecretProvenance: () => ({ complete: true, columns: {} }), createTableRowSecretProvenanceFromRegistry: () => ({ complete: true, columns: {} }), - loadTableRowSecretProvenance: mocks.loadTableRowSecretProvenance, + TableRowProvenanceReader: class { + constructor(scope: unknown, selectedColumnIds: unknown) { + mocks.createProvenanceReader(scope, selectedColumnIds) + } + exportProvenance = mocks.exportProvenance + }, })) vi.mock('@/executor/utils/resolved-secret-trace-registry', () => ({ ResolvedSecretTraceRegistry: class { @@ -151,6 +165,11 @@ describe('the workflow half of a table cell', () => { beforeEach(() => { vi.clearAllMocks() resetDbChainMock() + mocks.getRowSummaryById.mockImplementation((tableId, rowId, workspaceId) => + mocks.getRowById(tableId, rowId, workspaceId) + ) + mocks.flattenWorkflowOutputs.mockReturnValue([]) + mocks.normalizeInputFormatValue.mockReturnValue([]) mocks.getTableById.mockResolvedValue(TABLE) mocks.getRowById.mockResolvedValue({ id: 'row-1', @@ -163,7 +182,7 @@ describe('the workflow half of a table cell', () => { mocks.markWorkflowGroupPickedUp.mockResolvedValue('picked-up') mocks.loadDeployedWorkflowState.mockResolvedValue({ blocks: {}, edges: [] }) mocks.findStartBlock.mockReturnValue({ blockId: 'start-1', block: { subBlocks: {} } }) - mocks.loadTableRowSecretProvenance.mockResolvedValue({ + mocks.exportProvenance.mockReturnValue({ scope: { userId: 'workflow-owner', workspaceId: 'workspace-1' }, byRowId: {}, }) @@ -212,6 +231,53 @@ describe('the workflow half of a table cell', () => { expect(options.billingAttribution).toBe(BILLING) }, 20_000) + it.each([false, true])( + 'captures fresh workflow inputs and explicit output mappings (%s)', + async (mapOutput) => { + const group = { + ...GROUP, + outputs: [{ columnName: 'col-output', blockId: 'agent', path: 'content' }], + inputMappings: mapOutput ? [{ columnName: 'col-output', inputName: 'explicit' }] : [], + } + mocks.flattenWorkflowOutputs.mockReturnValue([{ blockId: 'agent', path: 'content' }]) + mocks.normalizeInputFormatValue.mockReturnValue([{ name: 'explicit' }]) + mocks.getTableById.mockResolvedValue({ + ...TABLE, + schema: { + columns: [ + { id: 'col-input', name: 'Input', type: 'string' }, + { id: 'col-output', name: 'Output', type: 'string' }, + ], + workflowGroups: [group], + }, + }) + mocks.getRowSummaryById.mockResolvedValue({ + id: 'row-1', + data: { 'col-input': 'current-value', 'col-output': 'output-secret' }, + updatedAt: new Date('2026-09-14T00:00:00Z'), + }) + + await runRowCascadeLoop(PAYLOAD) + + expect(mocks.createProvenanceReader).toHaveBeenCalledExactlyOnceWith( + { userId: 'workflow-owner', workspaceId: 'workspace-1' }, + new Set(mapOutput ? ['col-input', 'col-output'] : ['col-input']) + ) + const input = mocks.executeWorkflow.mock.calls[0][2] + expect(input.row).toEqual({ Input: 'current-value' }) + expect(input.rawRow).toEqual(input.row) + expect(input.headers).toEqual(['Input']) + if (mapOutput) expect(input.explicit).toBe('output-secret') + else expect(input.Input).toBe('current-value') + expect(mocks.markWorkflowGroupPickedUp.mock.invocationCallOrder[0]).toBeLessThan( + mocks.getRowSummaryById.mock.invocationCallOrder[0] + ) + expect(mocks.getRowSummaryById.mock.invocationCallOrder[0]).toBeLessThan( + mocks.executeWorkflow.mock.invocationCallOrder[0] + ) + } + ) + it('declares an explicit null for an actorless auto-fire', async () => { await runRowCascadeLoop({ ...PAYLOAD, capabilityGovernedUserId: null }) diff --git a/apps/sim/executor/utils/resolved-secret-trace-registry.test.ts b/apps/sim/executor/utils/resolved-secret-trace-registry.test.ts index 785b209fdc0..333b3100eac 100644 --- a/apps/sim/executor/utils/resolved-secret-trace-registry.test.ts +++ b/apps/sim/executor/utils/resolved-secret-trace-registry.test.ts @@ -2347,3 +2347,117 @@ describe('unredacted catalog exemption', () => { expect(registry.getUnredactedSecretNames()).toEqual([]) }) }) + +describe('current environment resolutions', () => { + const scope = { userId: 'user-1', workspaceId: 'workspace-1' } + const oldEntry = { + name: 'API_KEY', + plaintext: 'old-personal-secret', + encryptedValue: 'encrypted-old', + scope: 'personal' as const, + ownerUserId: scope.userId, + } + const environment = { + personalEncrypted: { API_KEY: oldEntry.encryptedValue }, + personalDecrypted: { API_KEY: oldEntry.plaintext }, + personalOwners: { API_KEY: scope.userId }, + workspaceEncrypted: { API_KEY: 'encrypted-current', UNRELATED: 'encrypted-unrelated' }, + workspaceDecrypted: { API_KEY: 'current-workspace-secret', UNRELATED: 'unrelated-secret' }, + scope, + } + + it('uses the current workspace value while preserving earlier active values and sibling snapshots', () => { + const parent = new ResolvedSecretTraceRegistry([oldEntry], scope) + parent.recordResolved(oldEntry.name, oldEntry.plaintext, { propagated: true }) + const sibling = parent.forkForInputPaths([]) + const current = parent.forkForToolCall() + + expect( + current.recordResolvedFromEnvironment('API_KEY', 'current-workspace-secret', environment, { + path: ['apiKey'], + propagated: true, + }) + ).toBe(true) + expect(current.isComplete()).toBe(true) + expect(current.getActiveMatches()).toEqual( + expect.arrayContaining([ + { plaintext: oldEntry.plaintext, replacement: '{{API_KEY}}' }, + { plaintext: 'current-workspace-secret', replacement: '{{API_KEY}}' }, + ]) + ) + expect(current.getResolvedSecretUsage()).toEqual( + expect.arrayContaining([ + { name: 'API_KEY', scope: 'personal', ownerUserId: scope.userId }, + { name: 'API_KEY', scope: 'workspace', ownerUserId: null }, + ]) + ) + expect(sibling.recordResolved('API_KEY', oldEntry.plaintext)).toBe(true) + expect(sibling.isComplete()).toBe(true) + expect(parent.getActiveMatches()).toEqual([ + { plaintext: oldEntry.plaintext, replacement: '{{API_KEY}}' }, + ]) + parent.mergeToolCallRegistry(current) + expect(parent.getActiveMatches()).toEqual(current.getActiveMatches()) + expect(current.recordResolved('UNRELATED', 'unrelated-secret')).toBe(false) + }) + + it('rejects a removed secret instead of vouching from the old catalog', () => { + const registry = new ResolvedSecretTraceRegistry([oldEntry], scope) + registry.recordResolved(oldEntry.name, oldEntry.plaintext, { propagated: true }) + + expect( + registry.recordResolvedFromEnvironment(oldEntry.name, oldEntry.plaintext, { + personalEncrypted: {}, + personalDecrypted: {}, + workspaceEncrypted: {}, + workspaceDecrypted: {}, + scope, + }) + ).toBe(false) + expect(registry.isComplete()).toBe(false) + expect(registry.getActiveMatches()).toEqual([ + { plaintext: oldEntry.plaintext, replacement: '{{API_KEY}}' }, + ]) + }) + + it.each([ + { userId: 'another-user', workspaceId: scope.workspaceId }, + { userId: scope.userId, workspaceId: 'another-workspace' }, + ])('rejects a snapshot from a different scope: %j', (otherScope) => { + const registry = new ResolvedSecretTraceRegistry([], scope) + + expect( + registry.recordResolvedFromEnvironment('API_KEY', 'current-workspace-secret', { + ...environment, + scope: otherScope, + }) + ).toBe(false) + expect(registry.isComplete()).toBe(false) + expect(registry.getActiveMatches()).toEqual([]) + }) + + it('does not trust plaintext that differs from the current snapshot', () => { + const registry = new ResolvedSecretTraceRegistry([oldEntry], scope) + + expect(registry.recordResolvedFromEnvironment('API_KEY', oldEntry.plaintext, environment)).toBe( + false + ) + expect(registry.isComplete()).toBe(false) + expect(registry.getActiveMatches()).toEqual([]) + }) + + it('preserves the current workspace unredacted setting without changing a personal contribution', () => { + const registry = new ResolvedSecretTraceRegistry([oldEntry], scope) + registry.recordResolved(oldEntry.name, oldEntry.plaintext, { propagated: true }) + + expect( + registry.recordResolvedFromEnvironment('API_KEY', 'current-workspace-secret', { + ...environment, + workspaceUnredactedKeys: ['API_KEY'], + }) + ).toBe(true) + expect(registry.getActiveMatches()).toEqual([ + { plaintext: oldEntry.plaintext, replacement: '{{API_KEY}}' }, + ]) + }) +}) diff --git a/apps/sim/executor/utils/resolved-secret-trace-registry.ts b/apps/sim/executor/utils/resolved-secret-trace-registry.ts index 432e0178a85..925a6d18bf5 100644 --- a/apps/sim/executor/utils/resolved-secret-trace-registry.ts +++ b/apps/sim/executor/utils/resolved-secret-trace-registry.ts @@ -1175,6 +1175,46 @@ export class ResolvedSecretTraceRegistry { child.copyResolvedInputPathsTo(this) } + /** + * Vouches for one explicit resolution using the same authorized environment snapshot that + * supplied its value. A long-lived Copilot turn can outlive a secret addition or rotation; + * refreshing this name leaves earlier active values and sibling call registries intact. + */ + recordResolvedFromEnvironment( + name: string, + resolvedValue: string, + environment: CreateResolvedSecretTraceRegistryOptions, + options: { path?: ResolvedSecretInputPath; propagated?: boolean } = {} + ): boolean { + if (!scopesMatch(this.scope, environment.scope)) { + this.markIncomplete('tool-call-scope-mismatch') + return false + } + + const encryptedValue = hasOwn(environment.workspaceEncrypted, name) + ? environment.workspaceEncrypted[name] + : environment.personalEncrypted[name] + const entry = + typeof encryptedValue === 'string' && encryptedValue.length > 0 + ? buildEffectiveCatalogEntry( + environment, + new Set(environment.decryptionFailures ?? []), + new Set(environment.workspaceUnredactedKeys ?? []), + name, + encryptedValue + ) + : undefined + if (!entry || entry.plaintext !== resolvedValue || !this.addCatalogEntry(entry)) { + if (options.path?.length) { + this.markInputPathIncomplete(options.path, 'unverified-resolved-entry') + } else { + this.markIncomplete('unverified-resolved-entry') + } + return false + } + return this.recordResolvedAtInputPath(name, resolvedValue, options.path, options) + } + /** Activates a configured secret only when the resolved runtime value matches its catalog value. */ recordResolved( name: string, diff --git a/apps/sim/lib/copilot/tools/handlers/upload-file-reader.test.ts b/apps/sim/lib/copilot/tools/handlers/upload-file-reader.test.ts index bb2deb7b37c..14e93a3ba3c 100644 --- a/apps/sim/lib/copilot/tools/handlers/upload-file-reader.test.ts +++ b/apps/sim/lib/copilot/tools/handlers/upload-file-reader.test.ts @@ -23,13 +23,15 @@ vi.mock('@/lib/uploads/contexts/workspace/workspace-file-manager', () => ({ /** A buffer beginning with the ZIP local-file-header magic (PK\x03\x04). */ const ZIP_SHAPED = Buffer.from([0x50, 0x4b, 0x03, 0x04, 0x14, 0x00, 0x00, 0x00]) -import { WorkspaceFileGrepError } from '@/lib/copilot/vfs/operations' import { findMothershipUploadRowByChatAndName, grepChatUpload, + grepChatUploadWithProvenance, listChatUploads, readChatUpload, -} from './upload-file-reader' + readChatUploadWithProvenance, +} from '@/lib/copilot/tools/handlers/upload-file-reader' +import { WorkspaceFileGrepError } from '@/lib/copilot/vfs/operations' const CHAT_ID = '11111111-1111-1111-1111-111111111111' const NOW = new Date('2026-05-05T00:00:00.000Z') @@ -49,6 +51,7 @@ function makeRow(overrides: Partial> = {}) { deletedAt: null, uploadedAt: NOW, updatedAt: NOW, + contentUpdatedAt: NOW, ...overrides, } } @@ -164,6 +167,24 @@ describe('readChatUpload', () => { mockReadFileRecord.mockReset() }) + it('captures the content revision before read and grep can race with upload promotion', async () => { + const row = makeRow({ displayName: 'note.txt', contentType: 'text/plain' }) + const result = { content: 'note', totalLines: 1 } + for (const read of [ + () => readChatUploadWithProvenance('note.txt', CHAT_ID), + () => grepChatUploadWithProvenance('note.txt', CHAT_ID, 'note'), + ]) { + mockOrderByThenLimit([row]) + mockReadFileRecord.mockResolvedValueOnce(result) + expect((await read())?.file).toEqual({ + fileId: row.id, + key: row.key, + context: 'mothership', + contentUpdatedAt: NOW, + }) + } + }) + it('reads the row resolved by the suffixed displayName', async () => { const row = makeRow({ id: 'wf_2', displayName: 'image (2).png' }) mockOrderByThenLimit([row]) diff --git a/apps/sim/lib/copilot/tools/handlers/upload-file-reader.ts b/apps/sim/lib/copilot/tools/handlers/upload-file-reader.ts index b6546bc54d6..e09f1736cbd 100644 --- a/apps/sim/lib/copilot/tools/handlers/upload-file-reader.ts +++ b/apps/sim/lib/copilot/tools/handlers/upload-file-reader.ts @@ -95,6 +95,7 @@ function toWorkspaceFileRecord(row: WorkspaceFileRow): WorkspaceFileRecord { deletedAt: row.deletedAt, uploadedAt: row.uploadedAt, updatedAt: row.updatedAt, + contentUpdatedAt: row.contentUpdatedAt, storageContext: 'mothership', } } @@ -207,7 +208,12 @@ export async function readChatUploadWithProvenance( if (!result) return null return { value: result, - file: { fileId: record.id, key: record.key, context: 'mothership' }, + file: { + fileId: record.id, + key: record.key, + context: 'mothership', + contentUpdatedAt: row.contentUpdatedAt, + }, view: isReadableFileType(record.type) ? 'complete' : 'derived', } } catch (err) { @@ -260,6 +266,11 @@ export async function grepChatUploadWithProvenance( const uploadsPath = `uploads/${canonicalUploadKey(record.name)}` return { value: grepReadResult(uploadsPath, result, pattern, uploadsPath, options), - file: { fileId: record.id, key: record.key, context: 'mothership' }, + file: { + fileId: record.id, + key: record.key, + context: 'mothership', + contentUpdatedAt: row.contentUpdatedAt, + }, } } diff --git a/apps/sim/lib/copilot/tools/server/env-reference.test.ts b/apps/sim/lib/copilot/tools/server/env-reference.test.ts new file mode 100644 index 00000000000..461bc7f125c --- /dev/null +++ b/apps/sim/lib/copilot/tools/server/env-reference.test.ts @@ -0,0 +1,69 @@ +/** + * @vitest-environment node + */ +import { environmentUtilsMockFns, resetEnvironmentUtilsMock } from '@sim/testing' +import { afterEach, describe, expect, it } from 'vitest' +import { resolveEnvReferenceSecretArg } from '@/lib/copilot/tools/server/env-reference' +import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' + +const scope = { userId: 'user-1', workspaceId: 'workspace-1' } + +describe('resolveEnvReferenceSecretArg', () => { + afterEach(resetEnvironmentUtilsMock) + + it('protects a password created after the turn registry was initialized', async () => { + const registry = new ResolvedSecretTraceRegistry([], scope) + environmentUtilsMockFns.mockGetEffectiveEnvironmentSnapshot.mockResolvedValueOnce({ + personalEncrypted: {}, + personalDecrypted: {}, + workspaceEncrypted: { CHAT_PASSWORD: 'encrypted-password' }, + workspaceDecrypted: { CHAT_PASSWORD: 'newly-created-password' }, + personalOwners: {}, + conflicts: [], + decryptionFailures: [], + workspaceUnredactedKeys: [], + }) + + expect( + await resolveEnvReferenceSecretArg({ + ...scope, + value: '{{CHAT_PASSWORD}}', + argName: 'password', + registry, + }) + ).toEqual({ value: 'newly-created-password' }) + expect(registry.isComplete()).toBe(true) + expect(registry.getActiveMatches()).toEqual([ + { plaintext: 'newly-created-password', replacement: '{{CHAT_PASSWORD}}' }, + ]) + }) + + it('does not resolve a removed password from the previous catalog', async () => { + const registry = new ResolvedSecretTraceRegistry( + [{ name: 'CHAT_PASSWORD', plaintext: 'old-password', encryptedValue: 'encrypted-old' }], + scope + ) + + const result = await resolveEnvReferenceSecretArg({ + ...scope, + value: '{{CHAT_PASSWORD}}', + argName: 'password', + registry, + }) + + expect(result).not.toHaveProperty('value') + expect(result.error).toContain('not set') + expect(registry.getActiveMatches()).toEqual([]) + }) + + it('leaves literal passwords alone without loading the environment', async () => { + expect( + await resolveEnvReferenceSecretArg({ + ...scope, + value: '$literal_password', + argName: 'password', + }) + ).toEqual({ value: '$literal_password' }) + expect(environmentUtilsMockFns.mockGetEffectiveEnvironmentSnapshot).not.toHaveBeenCalled() + }) +}) diff --git a/apps/sim/lib/copilot/tools/server/env-reference.ts b/apps/sim/lib/copilot/tools/server/env-reference.ts index e2339a22ab5..94e7bf9da66 100644 --- a/apps/sim/lib/copilot/tools/server/env-reference.ts +++ b/apps/sim/lib/copilot/tools/server/env-reference.ts @@ -1,4 +1,4 @@ -import { getEffectiveDecryptedEnv } from '@/lib/environment/utils' +import { getEffectiveEnvironmentSnapshot } from '@/lib/environment/utils' import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' /** @@ -26,14 +26,17 @@ export async function resolveEnvReferenceSecretArg(args: { const braced = value.match(/^\{\{\s*([A-Za-z_][A-Za-z0-9_]*)\s*\}\}$/) if (!braced) return { value } const name = braced[1] - const env = await getEffectiveDecryptedEnv(args.userId, args.workspaceId) + const environment = await getEffectiveEnvironmentSnapshot(args.userId, args.workspaceId) + const env = { ...environment.personalDecrypted, ...environment.workspaceDecrypted } const resolved = env[name] if (resolved === undefined || resolved === '') { return { error: `Environment variable "${name}" referenced by ${args.argName} is not set for this workspace or user. Set it first, or pass the raw value.`, } } - // Activate on the call's egress registry so an accidental echo is redacted. - args.registry?.recordResolved(name, resolved) + args.registry?.recordResolvedFromEnvironment(name, resolved, { + ...environment, + scope: { userId: args.userId, workspaceId: args.workspaceId }, + }) return { value: resolved } } diff --git a/apps/sim/lib/copilot/tools/server/knowledge/knowledge-base.test.ts b/apps/sim/lib/copilot/tools/server/knowledge/knowledge-base.test.ts index 298d5d7f270..4e0977f410a 100644 --- a/apps/sim/lib/copilot/tools/server/knowledge/knowledge-base.test.ts +++ b/apps/sim/lib/copilot/tools/server/knowledge/knowledge-base.test.ts @@ -89,11 +89,11 @@ vi.mock('@/lib/copilot/chat/organization-chats', () => ({ vi.mock('@/lib/copilot/generated/tool-catalog-v1', () => ({ ManageKnowledgeBase: { id: 'manage_knowledge_base' }, })) -const { mockGetEffectiveDecryptedEnv } = vi.hoisted(() => ({ - mockGetEffectiveDecryptedEnv: vi.fn(), +const { mockGetEffectiveEnvironmentSnapshot } = vi.hoisted(() => ({ + mockGetEffectiveEnvironmentSnapshot: vi.fn(), })) vi.mock('@/lib/environment/utils', () => ({ - getEffectiveDecryptedEnv: mockGetEffectiveDecryptedEnv, + getEffectiveEnvironmentSnapshot: mockGetEffectiveEnvironmentSnapshot, })) vi.mock('@/lib/core/telemetry', () => ({ PlatformEvents: { @@ -809,7 +809,12 @@ describe('manage_knowledge_base trusted application delegation', () => { it.each(['{{SIM_GITHUB_PAT}}', '$SIM_GITHUB_PAT', 'SIM_GITHUB_PAT'])( 'resolves the %s environment reference into the connector API key', async (ref) => { - mockGetEffectiveDecryptedEnv.mockResolvedValue({ SIM_GITHUB_PAT: 'ghp_realtoken' }) + mockGetEffectiveEnvironmentSnapshot.mockResolvedValue({ + personalEncrypted: {}, + personalDecrypted: {}, + workspaceEncrypted: { SIM_GITHUB_PAT: 'encrypted-token' }, + workspaceDecrypted: { SIM_GITHUB_PAT: 'ghp_realtoken' }, + }) const result = await knowledgeBaseServerTool.execute( { @@ -828,7 +833,12 @@ describe('manage_knowledge_base trusted application delegation', () => { ) it('names the missing variable instead of sending a placeholder upstream', async () => { - mockGetEffectiveDecryptedEnv.mockResolvedValue({}) + mockGetEffectiveEnvironmentSnapshot.mockResolvedValue({ + personalEncrypted: {}, + personalDecrypted: {}, + workspaceEncrypted: {}, + workspaceDecrypted: {}, + }) const result = await knowledgeBaseServerTool.execute( { @@ -849,7 +859,12 @@ describe('manage_knowledge_base trusted application delegation', () => { }) it('passes a raw API key through untouched', async () => { - mockGetEffectiveDecryptedEnv.mockResolvedValue({ SIM_GITHUB_PAT: 'ghp_realtoken' }) + mockGetEffectiveEnvironmentSnapshot.mockResolvedValue({ + personalEncrypted: {}, + personalDecrypted: {}, + workspaceEncrypted: { SIM_GITHUB_PAT: 'encrypted-token' }, + workspaceDecrypted: { SIM_GITHUB_PAT: 'ghp_realtoken' }, + }) const result = await knowledgeBaseServerTool.execute( { diff --git a/apps/sim/lib/copilot/tools/server/knowledge/knowledge-base.ts b/apps/sim/lib/copilot/tools/server/knowledge/knowledge-base.ts index a610f02b6b2..bcd994f9402 100644 --- a/apps/sim/lib/copilot/tools/server/knowledge/knowledge-base.ts +++ b/apps/sim/lib/copilot/tools/server/knowledge/knowledge-base.ts @@ -19,7 +19,7 @@ import { } from '@/lib/copilot/tools/server/base-tool' import { asOrchestrationError } from '@/lib/core/orchestration/types' import { PlatformEvents } from '@/lib/core/telemetry' -import { getEffectiveDecryptedEnv } from '@/lib/environment/utils' +import { getEffectiveEnvironmentSnapshot } from '@/lib/environment/utils' import { isKnowledgeMemberAccessAvailable } from '@/lib/knowledge/access/availability' import { addWorkspaceFilesToKnowledgeBase } from '@/lib/knowledge/application/add-workspace-files' import { KnowledgeUsageLimitExceededError } from '@/lib/knowledge/application/billing' @@ -96,7 +96,8 @@ async function resolveConnectorApiKey( const braced = apiKey.match(/^\{\{\s*([A-Za-z_][A-Za-z0-9_]*)\s*\}\}$/) const dollar = apiKey.match(/^\$([A-Za-z_][A-Za-z0-9_]*)$/) const referencedName = braced?.[1] ?? dollar?.[1] - const env = await getEffectiveDecryptedEnv(context.userId, workspaceId) + const environment = await getEffectiveEnvironmentSnapshot(context.userId, workspaceId) + const env = { ...environment.personalDecrypted, ...environment.workspaceDecrypted } const name = referencedName ?? (Object.hasOwn(env, apiKey) ? apiKey : undefined) if (!name) return { apiKey } const value = env[name] @@ -105,9 +106,10 @@ async function resolveConnectorApiKey( error: `Environment variable "${name}" is not set for this workspace or user, so it cannot be used as the connector API key. Set it first, pass a different {{ENV_VAR}} reference, or pass the raw key.`, } } - // Activate the resolved secret on the call's egress registry so any - // accidental echo of it (provider error bodies, logs) is redacted. - context.resolvedSecretTraceRegistry?.recordResolved(name, value) + context.resolvedSecretTraceRegistry?.recordResolvedFromEnvironment(name, value, { + ...environment, + scope: { userId: context.userId, workspaceId }, + }) return { apiKey: value } } diff --git a/apps/sim/lib/copilot/tools/server/table/user-table.test.ts b/apps/sim/lib/copilot/tools/server/table/user-table.test.ts index 3b510d3d26d..2d247f6847d 100644 --- a/apps/sim/lib/copilot/tools/server/table/user-table.test.ts +++ b/apps/sim/lib/copilot/tools/server/table/user-table.test.ts @@ -34,7 +34,8 @@ const { mockExecuteCopilotFileUseCase, mockExecuteCopilotWorkflowUseCase, mockLoadWorkspaceFileContext, - mockLoadTableRowSecretProvenance, + mockCreateTableRowProvenanceReader, + mockExportTableRowProvenance, mockResolveWorkflowContext, fakeEnrichment, } = vi.hoisted(() => ({ @@ -65,7 +66,8 @@ const { mockExecuteCopilotFileUseCase: vi.fn(), mockExecuteCopilotWorkflowUseCase: vi.fn(), mockLoadWorkspaceFileContext: vi.fn(), - mockLoadTableRowSecretProvenance: vi.fn(), + mockCreateTableRowProvenanceReader: vi.fn(), + mockExportTableRowProvenance: vi.fn(), mockResolveWorkflowContext: vi.fn(), fakeEnrichment: { id: 'work-email', @@ -229,7 +231,12 @@ vi.mock('@/lib/table/rows/secret-provenance', () => ({ Object.keys(data).map((columnId) => [columnId, { version: 1, complete: true, entries: [] }]) ), }), - loadTableRowSecretProvenance: mockLoadTableRowSecretProvenance, + TableRowProvenanceReader: class { + constructor(scope: unknown) { + mockCreateTableRowProvenanceReader(scope) + } + exportProvenance = mockExportTableRowProvenance + }, })) vi.mock('@/lib/table/jobs/service', () => ({ @@ -271,7 +278,7 @@ import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-tr beforeEach(() => { mockLoadWorkspaceFileContext.mockResolvedValue({ workspaceId: 'workspace-1' }) - mockLoadTableRowSecretProvenance.mockResolvedValue({ + mockExportTableRowProvenance.mockReturnValue({ version: 1, complete: true, entries: [], @@ -1223,9 +1230,12 @@ describe('userTableServerTool.query_rows', () => { ) expect(result.success).toBe(true) - expect(mockLoadTableRowSecretProvenance).toHaveBeenCalledWith( - expect.arrayContaining([expect.objectContaining({ id: 'row_1' })]), - { userId: 'user-1', workspaceId: 'workspace-1' } + expect(mockCreateTableRowProvenanceReader).toHaveBeenCalledWith({ + userId: 'user-1', + workspaceId: 'workspace-1', + }) + expect(mockQueryRows.mock.calls[0][3]).toEqual( + expect.objectContaining({ exportProvenance: mockExportTableRowProvenance }) ) expect(importProvenance).toHaveBeenCalledWith( expect.objectContaining({ diff --git a/apps/sim/lib/internal/file/parser.test.ts b/apps/sim/lib/internal/file/parser.test.ts index c3018aeef8c..9a6fd2a5e68 100644 --- a/apps/sim/lib/internal/file/parser.test.ts +++ b/apps/sim/lib/internal/file/parser.test.ts @@ -92,6 +92,8 @@ const { provider: 's3', bucket: 'sim-execution-files', containerName: 'execution-files', + workspaceBucket: 'sim-workspace-files', + workspaceContainerName: 'workspace-files', }, mockGetBlobContainerClient: vi.fn(), } @@ -132,7 +134,13 @@ vi.mock('@/lib/uploads', () => ({ })) vi.mock('@/lib/uploads/config', () => ({ - getStorageConfig: () => storageConfig, + getStorageConfig: (context: string) => + context === 'workspace' + ? { + bucket: storageConfig.workspaceBucket, + containerName: storageConfig.workspaceContainerName, + } + : storageConfig, S3_CONFIG: {}, get USE_S3_STORAGE() { return storageConfig.provider === 's3' @@ -502,6 +510,143 @@ describe('file parser operation', () => { expect(inputValidationMockFns.mockSecureFetchWithPinnedIP).not.toHaveBeenCalled() }) + it.each([ + { + provider: 's3', + baseUrl: 'https://sim-workspace-files.s3.us-east-1.amazonaws.com', + }, + { + provider: 's3', + baseUrl: 'https://s3.us-east-1.amazonaws.com/sim-workspace-files', + }, + { + provider: 'blob', + baseUrl: 'https://exampleaccount.blob.core.windows.net/workspace-files', + }, + { + provider: 'gcs', + baseUrl: 'https://sim-workspace-files.storage.googleapis.com', + }, + { + provider: 'gcs', + baseUrl: 'https://storage.googleapis.com/sim-workspace-files', + }, + ])('preserves owned workspace URL lineage for $baseUrl', async ({ provider, baseUrl }) => { + storageConfig.provider = provider + mockGetBlobContainerClient.mockImplementation((containerName: string) => ({ + url: `https://exampleaccount.blob.core.windows.net/${containerName}`, + })) + const key = 'workspace/workspace-id/report.txt' + const source = { + identity: { fileId: 'canonical-file', key, context: 'workspace' }, + ownerUserId: 'test-user-id', + } + const entries = [ + { name: 'SECRET', encryptedValue: 'encrypted-value', sourceUserId: 'test-user-id' }, + ] + mockResolveProvenanceSource.mockResolvedValue(source) + mockGetBoundProvenance.mockResolvedValue({ status: 'exact', entries }) + + const response = await POST( + createMockRequest( + 'POST', + { filePath: `${baseUrl}/${key}?signature=test` }, + { 'x-sim-request-private-tool-metadata': 'resolved-secret-provenance-v1' } + ) + ) + + expect((await response.json()).success).toBe(true) + expect(mockResolveProvenanceSource).toHaveBeenCalledWith( + { key, context: 'workspace' }, + expect.objectContaining({ workspaceId: 'workspace-id' }) + ) + expect(mockUploadExecutionFile.mock.calls[0][5]).toEqual({ status: 'exact', entries }) + expect(mockGetFileContentProvenance).toHaveBeenCalledWith( + expect.any(Object), + 'workspace-id', + [source], + expect.any(AbortSignal) + ) + expect(inputValidationMockFns.mockSecureFetchWithPinnedIP).not.toHaveBeenCalled() + }) + + it('does not reclassify a rejected owned URL as external ingress', async () => { + mockResolveProvenanceSource.mockRejectedValue(new Error('File not found')) + + const response = await POST( + createMockRequest('POST', { + filePath: + 'https://sim-workspace-files.s3.us-east-1.amazonaws.com/workspace/workspace-id/report.txt', + }) + ) + + expect((await response.json()).success).toBe(false) + expect(storageServiceMockFns.mockDownloadFile).not.toHaveBeenCalled() + expect(mockUploadExecutionFile).not.toHaveBeenCalled() + expect(inputValidationMockFns.mockSecureFetchWithPinnedIP).not.toHaveBeenCalled() + }) + + it('does not fall back to external ingress when configured storage resolution fails', async () => { + storageConfig.provider = 'blob' + mockGetBlobContainerClient.mockImplementation(() => { + throw new Error('Storage configuration unavailable') + }) + + const response = await POST( + createMockRequest('POST', { + filePath: + 'https://exampleaccount.blob.core.windows.net/workspace-files/workspace/workspace-id/report.txt', + }) + ) + + expect((await response.json()).success).toBe(false) + expect(inputValidationMockFns.mockSecureFetchWithPinnedIP).not.toHaveBeenCalled() + expect(mockUploadExecutionFile).not.toHaveBeenCalled() + }) + + it('records authenticated external binary downloads without a Sim-secret contribution', async () => { + const bytes = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x00]) + inputValidationMockFns.mockValidateUrlWithDNS.mockResolvedValue({ + isValid: true, + resolvedIP: '203.0.113.10', + }) + inputValidationMockFns.mockSecureFetchWithPinnedIP.mockResolvedValue( + new Response(bytes, { headers: { 'content-type': 'image/png' } }) + ) + permissionsMockFns.mockGetUserEntityPermissions.mockResolvedValue('write') + mockIsSupportedFileType.mockReturnValue(false) + + await POST( + createMockRequest('POST', { + filePath: 'https://files.slack.com/files-pri/T07-FAAA/download/image.png', + headers: { Authorization: 'Bearer xoxb-test-token' }, + }) + ) + + expect(inputValidationMockFns.mockSecureFetchWithPinnedIP).toHaveBeenCalledWith( + 'https://files.slack.com/files-pri/T07-FAAA/download/image.png', + '203.0.113.10', + expect.objectContaining({ headers: { Authorization: 'Bearer xoxb-test-token' } }) + ) + expect(mockUploadExecutionFile).toHaveBeenCalledWith( + expect.any(Object), + bytes, + 'image.png', + 'image/png', + 'test-user-id', + { status: 'exact', entries: [] } + ) + expect(mockUploadWorkspaceFile).toHaveBeenCalledWith( + 'workspace-id', + 'test-user-id', + bytes, + 'image.png', + 'image/png' + ) + expect(mockResolveProvenanceSource).not.toHaveBeenCalled() + expect(mockGetBoundProvenance).not.toHaveBeenCalled() + }) + it.each([ 'https://exampleaccount.blob.core.windows.net/execution-files', 'https://exampleaccount.blob.core.usgovcloudapi.net/execution-files', diff --git a/apps/sim/lib/internal/file/parser.ts b/apps/sim/lib/internal/file/parser.ts index 19f24e02ef7..470a1d7ebaf 100644 --- a/apps/sim/lib/internal/file/parser.ts +++ b/apps/sim/lib/internal/file/parser.ts @@ -563,7 +563,7 @@ function validateFilePath(filePath: string): { isValid: boolean; error?: string * so keying a cache by filename returns stale bytes. `fetchExternalUrlToWorkspace` * delegates to `uploadWorkspaceFile`, which suffix-disambiguates collisions on save. * - * URLs for our execution-files storage resolve through the authorized canonical + * URLs for our workspace and execution storage resolve through the authorized canonical * read path, keeping stored provenance bound to the same bytes the parser reads. */ async function handleExternalUrl( @@ -583,15 +583,14 @@ async function handleExternalUrl( const { getStorageConfig, S3_CONFIG, USE_S3_STORAGE, USE_BLOB_STORAGE, USE_GCS_STORAGE } = await import('@/lib/uploads/config') - const executionConfig = getStorageConfig('execution') - - let executionFileKey: string | undefined - try { - const parsedUrl = new URL(url) - - if (USE_S3_STORAGE && executionConfig.bucket) { + const parsedUrl = new URL(url) + let ownedFileKey: string | undefined + /** Workspace storage also holds mothership attachments. */ + for (const storageContext of ['execution', 'workspace'] as const) { + const storageConfig = getStorageConfig(storageContext) + if (USE_S3_STORAGE && storageConfig.bucket) { const endpointHost = S3_CONFIG.endpoint ? new URL(S3_CONFIG.endpoint).host : undefined - const bucketHostPrefix = `${executionConfig.bucket}.` + const bucketHostPrefix = `${storageConfig.bucket}.` const storageHost = parsedUrl.host.startsWith(bucketHostPrefix) ? parsedUrl.host.slice(bucketHostPrefix.length) : parsedUrl.host @@ -600,45 +599,42 @@ async function handleExternalUrl( : /^s3(?:[.-][a-z0-9-]+)?\.amazonaws\.com$/.test(storageHost) const bucketInHost = matchesStorageHost && parsedUrl.host.startsWith(bucketHostPrefix) const bucketInPath = - matchesStorageHost && parsedUrl.pathname.startsWith(`/${executionConfig.bucket}/`) + matchesStorageHost && parsedUrl.pathname.startsWith(`/${storageConfig.bucket}/`) if (bucketInHost || bucketInPath) { - executionFileKey = decodeURIComponent( + ownedFileKey = decodeURIComponent( bucketInHost ? parsedUrl.pathname.slice(1) - : parsedUrl.pathname.slice(executionConfig.bucket.length + 2) + : parsedUrl.pathname.slice(storageConfig.bucket.length + 2) ) } - } else if (USE_BLOB_STORAGE && executionConfig.containerName) { + } else if (USE_BLOB_STORAGE && storageConfig.containerName) { const { getBlobServiceClient } = await import('@/lib/uploads/providers/blob/client') const client = await getBlobServiceClient() - const containerUrl = new URL(client.getContainerClient(executionConfig.containerName).url) + const containerUrl = new URL(client.getContainerClient(storageConfig.containerName).url) const prefix = `${containerUrl.pathname.replace(/\/$/, '')}/` if (parsedUrl.origin === containerUrl.origin && parsedUrl.pathname.startsWith(prefix)) { - executionFileKey = decodeURIComponent(parsedUrl.pathname.slice(prefix.length)) + ownedFileKey = decodeURIComponent(parsedUrl.pathname.slice(prefix.length)) } - } else if (USE_GCS_STORAGE && executionConfig.bucket) { - const bucketInHost = - parsedUrl.hostname === `${executionConfig.bucket}.storage.googleapis.com` + } else if (USE_GCS_STORAGE && storageConfig.bucket) { + const bucketInHost = parsedUrl.hostname === `${storageConfig.bucket}.storage.googleapis.com` const bucketInPath = parsedUrl.hostname === 'storage.googleapis.com' && - parsedUrl.pathname.startsWith(`/${executionConfig.bucket}/`) + parsedUrl.pathname.startsWith(`/${storageConfig.bucket}/`) if (bucketInHost || bucketInPath) { - executionFileKey = decodeURIComponent( + ownedFileKey = decodeURIComponent( bucketInHost ? parsedUrl.pathname.slice(1) - : parsedUrl.pathname.slice(executionConfig.bucket.length + 2) + : parsedUrl.pathname.slice(storageConfig.bucket.length + 2) ) } } - } catch (error) { - logger.warn('Failed to parse URL for execution file check:', error) - executionFileKey = undefined + if (ownedFileKey) break } /** Read owned storage through its authorized, canonical bytes and provenance together. */ - if (executionFileKey) { + if (ownedFileKey) { return handleCloudFile( - executionFileKey, + ownedFileKey, fileType, userId, fileReadAccess, @@ -668,8 +664,13 @@ async function handleExternalUrl( let userFile: UserFile | undefined if (executionContext) { try { + /** + * External ingress carries no tracked Sim-secret contribution. Using a fetch credential + * does not classify the response bytes as derived secret content. + */ userFile = await uploadExecutionFile(executionContext, buffer, filename, mimeType, userId, { - status: 'unrecorded', + status: 'exact', + entries: [], }) logger.info(`Stored file in execution storage: ${filename}`, { key: userFile.key }) } catch (uploadError) { diff --git a/apps/sim/lib/knowledge/__integration__/external-file-provenance.integration.ts b/apps/sim/lib/knowledge/__integration__/external-file-provenance.integration.ts new file mode 100644 index 00000000000..0c0665585b9 --- /dev/null +++ b/apps/sim/lib/knowledge/__integration__/external-file-provenance.integration.ts @@ -0,0 +1,280 @@ +/** Real parser downloads, execution storage, durable sidecars, and enforced model admission. */ +import { mkdtempSync } from 'node:fs' +import { rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import path from 'node:path' +import { db } from '@sim/db' +import { + knowledgeBase, + organization, + user, + workspace, + workspaceFileColumns, + workspaceFiles, +} from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { isPlainRecord } from '@sim/utils/object' +import { eq, inArray } from 'drizzle-orm' +import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' + +const fixtureStorage = vi.hoisted(() => { + process.env.DURABLE_SECRET_PROVENANCE_ENFORCED_SURFACES = 'workspace-file' + return { root: '' } +}) +vi.mock('@/lib/uploads/core/setup.server', () => ({ + get UPLOAD_DIR_SERVER() { + return fixtureStorage.root + }, +})) + +import { fileParseBodySchema } from '@/lib/api/contracts/storage-transfer' +import { encryptSecret } from '@/lib/core/security/encryption' +import * as inputValidation from '@/lib/core/security/input-validation.server' +import { isUserFileWithMetadata } from '@/lib/core/utils/user-file' +import { isDurableSecretProvenanceEnforced } from '@/lib/execution/durable-secret-provenance-enforcement' +import { executeFileParserOperation } from '@/lib/internal/file/parser' +import { + createKnowledgeAclFixtureIds, + seedKnowledgeAclFixture, +} from '@/lib/knowledge/__integration__/seed-source-access-fixture' +import { uploadExecutionFile } from '@/lib/uploads/contexts/execution/execution-file-manager' +import { uploadWorkspaceFile } from '@/lib/uploads/contexts/workspace/workspace-file-manager' +import { + getBoundWorkspaceFileSecretProvenance, + importWorkspaceFileSecretProvenanceForModelView, + type WorkspaceFileSecretProvenance, + type WorkspaceFileSecretProvenanceIdentity, +} from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' +import { downloadFile } from '@/lib/uploads/core/storage-service' +import { createWorkspaceFileDelegatedPrincipal } from '@/lib/workspace-files/application/delegated-principal' +import type { UserFile } from '@/executor/types' +import { projectResolvedSecretModelContent } from '@/executor/utils/resolved-secret-content-projection' +import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' + +const fixtures: ReturnType[] = [] +const fetchSpy = vi.spyOn(inputValidation, 'secureFetchWithPinnedIP') +const validateUrlSpy = vi.spyOn(inputValidation, 'validateUrlWithDNS') +const EXTERNAL_IMAGE_URL = 'https://fixture.example.test/private/image.png' +const FIXTURE_SECRET = 'fixture-sim-secret-not-a-live-credential' +const IMAGE_BYTES = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x00, 0x01]) + +async function seed() { + const ids = createKnowledgeAclFixtureIds() + fixtures.push(ids) + await seedKnowledgeAclFixture(ids) + return { ...ids, workflowId: generateId(), executionId: generateId() } +} + +type Fixture = Awaited> + +async function parse(ids: Fixture, filePath: string, headers?: Record) { + const response = await executeFileParserOperation( + fileParseBodySchema.parse({ filePath, headers }), + { + principal: createWorkspaceFileDelegatedPrincipal({ + serviceId: 'executor', + subjectUserId: ids.aliceId, + workspaceId: ids.workspaceId, + delegationId: generateId(), + executionId: ids.executionId, + }), + workspaceId: ids.workspaceId, + workflowId: ids.workflowId, + executionId: ids.executionId, + attributedUserId: ids.aliceId, + fileAccessUserId: ids.aliceId, + } + ) + const body: unknown = await response.json() + expect(response.status, JSON.stringify(body)).toBe(200) + if ( + !isPlainRecord(body) || + body.success !== true || + !isPlainRecord(body.output) || + !isUserFileWithMetadata(body.output.file) || + typeof body.output.content !== 'string' + ) { + throw new Error(`Parser did not return an execution copy: ${JSON.stringify(body)}`) + } + return { file: body.output.file, content: body.output.content } +} + +async function identityFor(file: UserFile): Promise { + const [record] = await db + .select(workspaceFileColumns) + .from(workspaceFiles) + .where(eq(workspaceFiles.key, file.key)) + if (!record || record.context !== 'execution') { + throw new Error('Parser copy has no canonical execution metadata') + } + expect(record.secretProvenanceVersion).toBe(1) + return { + fileId: record.id, + key: record.key, + context: record.context, + contentUpdatedAt: record.contentUpdatedAt, + } +} + +beforeAll(() => { + fixtureStorage.root = mkdtempSync(path.join(tmpdir(), 'sim-external-file-provenance-')) + expect(isDurableSecretProvenanceEnforced('workspace-file')).toBe(true) +}) +beforeEach(() => { + fetchSpy.mockReset() + validateUrlSpy.mockReset() + validateUrlSpy.mockResolvedValue({ + isValid: true, + resolvedIP: '203.0.113.10', + originalHostname: new URL(EXTERNAL_IMAGE_URL).hostname, + }) + const response = new Response(IMAGE_BYTES) + fetchSpy.mockResolvedValue({ + ok: response.ok, + status: response.status, + statusText: response.statusText, + headers: new inputValidation.SecureFetchHeaders({ 'content-type': 'image/png' }), + body: response.body, + text: () => response.text(), + json: () => response.json(), + arrayBuffer: () => response.arrayBuffer(), + }) +}) +afterAll(async () => { + fetchSpy.mockRestore() + validateUrlSpy.mockRestore() + for (const ids of fixtures) { + await db.delete(knowledgeBase).where(eq(knowledgeBase.id, ids.knowledgeBaseId)) + await db.delete(workspace).where(eq(workspace.id, ids.workspaceId)) + await db.delete(organization).where(eq(organization.id, ids.organizationId)) + await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId])) + } + await rm(fixtureStorage.root, { recursive: true, force: true }) + await db.$client.end() +}) + +describe('external file ingress under durable enforcement', () => { + it('admits the exact downloaded binary bytes despite authenticated transport', async () => { + const ids = await seed() + const parsed = await parse(ids, EXTERNAL_IMAGE_URL, { + Authorization: 'Bearer fixture-download-credential', + }) + expect(fetchSpy).toHaveBeenCalledWith( + EXTERNAL_IMAGE_URL, + '203.0.113.10', + expect.objectContaining({ + headers: { Authorization: 'Bearer fixture-download-credential' }, + }) + ) + const identity = await identityFor(parsed.file) + expect(await downloadFile({ key: parsed.file.key, context: 'execution' })).toEqual(IMAGE_BYTES) + expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, identity)).toEqual({ + status: 'exact', + entries: [], + }) + expect( + await importWorkspaceFileSecretProvenanceForModelView({ + workspaceId: ids.workspaceId, + identity, + view: 'opaque', + }) + ).toBe(true) + + /** A real unrecorded control proves this suite has not silently disabled enforcement. */ + const unrecorded = await uploadExecutionFile( + ids, + IMAGE_BYTES, + 'unrecorded.png', + 'image/png', + ids.aliceId, + { status: 'unrecorded' } + ) + expect( + await importWorkspaceFileSecretProvenanceForModelView({ + workspaceId: ids.workspaceId, + identity: await identityFor(unrecorded), + view: 'opaque', + }) + ).toBe(false) + }) + + it('preserves and redacts owned text lineage while refusing opaque egress', async () => { + const ids = await seed() + const { encrypted } = await encryptSecret(FIXTURE_SECRET) + const provenance: WorkspaceFileSecretProvenance = { + status: 'exact', + entries: [ + { + name: 'FIXTURE_SECRET', + encryptedValue: encrypted, + sourceUserId: ids.aliceId, + sourceWorkspaceId: ids.workspaceId, + }, + ], + } + const content = `Saved secret: ${FIXTURE_SECRET}\n` + const source = await uploadWorkspaceFile( + ids.workspaceId, + ids.aliceId, + Buffer.from(content), + 'owned.txt', + 'text/plain', + { secretProvenance: provenance } + ) + const parsed = await parse(ids, source.url) + const identity = await identityFor(parsed.file) + expect(fetchSpy).not.toHaveBeenCalled() + expect(parsed.content).toBe(content) + expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, identity)).toEqual( + provenance + ) + expect( + await importWorkspaceFileSecretProvenanceForModelView({ + workspaceId: ids.workspaceId, + identity, + view: 'opaque', + }) + ).toBe(false) + + const registry = new ResolvedSecretTraceRegistry([], { + userId: ids.aliceId, + workspaceId: ids.workspaceId, + }) + expect( + await importWorkspaceFileSecretProvenanceForModelView({ + workspaceId: ids.workspaceId, + identity, + registry, + value: parsed.content, + view: 'complete', + }) + ).toBe(true) + const projected = projectResolvedSecretModelContent(parsed.content, registry) + expect(projected).toEqual({ safe: true, value: 'Saved secret: {{FIXTURE_SECRET}}\n' }) + }) + + it('keeps an owned unknown source unknown through its execution copy', async () => { + const ids = await seed() + const source = await uploadWorkspaceFile( + ids.workspaceId, + ids.aliceId, + IMAGE_BYTES, + 'unknown.png', + 'image/png', + { secretProvenance: { status: 'unknown' } } + ) + const parsed = await parse(ids, source.url) + const identity = await identityFor(parsed.file) + expect(fetchSpy).not.toHaveBeenCalled() + expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, identity)).toEqual({ + status: 'unknown', + }) + expect( + await importWorkspaceFileSecretProvenanceForModelView({ + workspaceId: ids.workspaceId, + identity, + view: 'opaque', + }) + ).toBe(false) + }) +}) diff --git a/apps/sim/lib/knowledge/__integration__/upload-read-provenance.integration.ts b/apps/sim/lib/knowledge/__integration__/upload-read-provenance.integration.ts new file mode 100644 index 00000000000..381c32a1a2b --- /dev/null +++ b/apps/sim/lib/knowledge/__integration__/upload-read-provenance.integration.ts @@ -0,0 +1,217 @@ +/** Real upload reads and save_upload must agree on the same bytes across promotion. */ +import { mkdtempSync } from 'node:fs' +import { rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import path from 'node:path' +import { db } from '@sim/db' +import { + copilotChats, + knowledgeBase, + organization, + user, + workspace, + workspaceFileColumns, + workspaceFiles, +} from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq, inArray } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' + +const fixtureStorage = vi.hoisted(() => ({ root: '' })) +vi.mock('@/lib/uploads/core/setup.server', () => ({ + get UPLOAD_DIR_SERVER() { + return fixtureStorage.root + }, +})) + +import { executeMaterializeFile } from '@/lib/copilot/tools/handlers/materialize-file' +import { + grepChatUploadWithProvenance, + readChatUploadWithProvenance, +} from '@/lib/copilot/tools/handlers/upload-file-reader' +import { + createKnowledgeAclFixtureIds, + seedKnowledgeAclFixture, +} from '@/lib/knowledge/__integration__/seed-source-access-fixture' +import { + generateWorkspaceFileKey, + trackChatUpload, +} from '@/lib/uploads/contexts/workspace/workspace-file-manager' +import { + getBoundWorkspaceFileSecretProvenance, + importWorkspaceFileSecretProvenanceForModelView, + replaceWorkspaceFileSecretProvenanceInTx, + type WorkspaceFileSecretProvenance, +} from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' +import { uploadFile } from '@/lib/uploads/core/storage-service' + +const fixtures: ReturnType[] = [] +const CONTENT = 'The same uploaded bytes remain readable after saving.\n' + +async function seedUpload(provenance?: WorkspaceFileSecretProvenance) { + const ids = createKnowledgeAclFixtureIds() + fixtures.push(ids) + await seedKnowledgeAclFixture(ids) + const chatId = generateId() + await db.insert(copilotChats).values({ + id: chatId, + userId: ids.aliceId, + workspaceId: ids.workspaceId, + type: 'mothership', + }) + const name = 'fixture.txt' + const key = generateWorkspaceFileKey(ids.workspaceId, name) + await uploadFile({ + file: Buffer.from(CONTENT), + fileName: name, + contentType: 'text/plain', + context: 'mothership', + customKey: key, + preserveKey: true, + metadata: { userId: ids.aliceId, workspaceId: ids.workspaceId, originalName: name }, + }) + await trackChatUpload( + ids.workspaceId, + ids.aliceId, + chatId, + key, + name, + 'text/plain', + CONTENT.length + ) + const [file] = await db + .select(workspaceFileColumns) + .from(workspaceFiles) + .where(eq(workspaceFiles.key, key)) + if (provenance) { + await db.transaction((tx) => + replaceWorkspaceFileSecretProvenanceInTx(tx, file.id, file.contentUpdatedAt, provenance) + ) + } + return { ...ids, chatId, file, name } +} + +async function promoteUpload(ids: Awaited>) { + const result = await executeMaterializeFile( + { fileNames: [ids.name], operation: 'save' }, + { + userId: ids.aliceId, + workspaceId: ids.workspaceId, + workflowId: '', + chatId: ids.chatId, + toolCallId: generateId(), + copilotToolExecution: true, + } + ) + expect(result.success).toBe(true) +} + +beforeAll(() => { + fixtureStorage.root = mkdtempSync(path.join(tmpdir(), 'sim-upload-read-provenance-')) +}) + +afterAll(async () => { + for (const ids of fixtures) { + await db.delete(knowledgeBase).where(eq(knowledgeBase.id, ids.knowledgeBaseId)) + await db.delete(workspace).where(eq(workspace.id, ids.workspaceId)) + await db.delete(organization).where(eq(organization.id, ids.organizationId)) + await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId])) + } + await rm(fixtureStorage.root, { recursive: true, force: true }) + await db.$client.end() +}) + +describe('chat upload reads racing with save_upload', () => { + it.each(['legacy', 'exact'] as const)( + 'keeps a %s upload read and grep valid across promotion', + async (kind) => { + const ids = await seedUpload(kind === 'exact' ? { status: 'exact', entries: [] } : undefined) + const read = await readChatUploadWithProvenance(ids.name, ids.chatId) + const grep = await grepChatUploadWithProvenance(ids.name, ids.chatId, 'uploaded') + expect(read?.value.content).toBe(CONTENT) + expect(read?.file).toBeDefined() + expect(grep.file).toBeDefined() + + /** Fixes the production interleaving: save commits after bytes are read, before admission. */ + await promoteUpload(ids) + + for (const envelope of [read, grep]) { + if (!envelope?.file) throw new Error('Upload read did not return its identity') + expect( + await importWorkspaceFileSecretProvenanceForModelView({ + workspaceId: ids.workspaceId, + actorUserId: ids.aliceId, + identity: envelope.file, + value: envelope.value, + view: 'opaque', + }) + ).toBe(true) + expect(envelope.file.contentUpdatedAt).toEqual(ids.file.contentUpdatedAt) + } + expect(await readChatUploadWithProvenance(ids.name, ids.chatId)).toBeNull() + } + ) + + it('still rejects changed bytes, mismatched scope or key, and a promotion without a captured revision', async () => { + const ids = await seedUpload({ status: 'exact', entries: [] }) + const read = await readChatUploadWithProvenance(ids.name, ids.chatId) + if (!read?.file) throw new Error('Upload read did not return its identity') + await promoteUpload(ids) + for (const identity of [ + { ...read.file, key: `${read.file.key}-different` }, + { ...read.file, contentUpdatedAt: undefined }, + { ...read.file, context: 'execution' as const }, + ]) { + expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, identity)).toEqual({ + status: 'unknown', + }) + } + expect(await getBoundWorkspaceFileSecretProvenance(generateId(), read.file)).toEqual({ + status: 'unknown', + }) + const nextRevision = new Date(ids.file.contentUpdatedAt.getTime() + 1) + await db.transaction(async (tx) => { + await tx + .update(workspaceFiles) + .set({ contentUpdatedAt: nextRevision }) + .where(eq(workspaceFiles.id, ids.file.id)) + await replaceWorkspaceFileSecretProvenanceInTx(tx, ids.file.id, nextRevision, { + status: 'exact', + entries: [], + }) + }) + expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, read.file)).toEqual({ + status: 'unknown', + }) + }) + + it.each([ + { status: 'unknown' }, + { + status: 'exact', + entries: [ + { name: 'API_KEY', encryptedValue: 'fixture-ciphertext', sourceUserId: 'fixture-author' }, + ], + }, + ])( + 'does not relax an opaque read with %j provenance when the upload is promoted', + async (provenance) => { + const ids = await seedUpload(provenance) + const read = await readChatUploadWithProvenance(ids.name, ids.chatId) + if (!read?.file) throw new Error('Upload read did not return its identity') + await promoteUpload(ids) + expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, read.file)).toEqual( + provenance + ) + expect( + await importWorkspaceFileSecretProvenanceForModelView({ + workspaceId: ids.workspaceId, + actorUserId: ids.aliceId, + identity: read.file, + value: read.value, + view: 'opaque', + }) + ).toBe(false) + } + ) +}) diff --git a/apps/sim/lib/table/__tests__/update-row.test.ts b/apps/sim/lib/table/__tests__/update-row.test.ts index f2027de5445..8370f7b192d 100644 --- a/apps/sim/lib/table/__tests__/update-row.test.ts +++ b/apps/sim/lib/table/__tests__/update-row.test.ts @@ -122,6 +122,9 @@ describe('updateRow — partial merge', () => { }) it('preserves columns not included in the partial update', async () => { + dbChainMockFns.returning.mockResolvedValueOnce([ + { ...EXISTING_ROW, data: { name: 'Alice', age: 31 }, updatedAt: PERSISTED_UPDATED_AT }, + ]) const result = await updateRow( { tableId: 'tbl-1', rowId: 'row-1', data: { age: 31 }, workspaceId: 'ws-1' }, TABLE, @@ -176,6 +179,9 @@ describe('updateRow — partial merge', () => { }) it('allows updating a single column without affecting others', async () => { + dbChainMockFns.returning.mockResolvedValueOnce([ + { ...EXISTING_ROW, data: { name: 'Bob', age: 30 }, updatedAt: PERSISTED_UPDATED_AT }, + ]) const result = await updateRow( { tableId: 'tbl-1', rowId: 'row-1', data: { name: 'Bob' }, workspaceId: 'ws-1' }, TABLE, @@ -187,6 +193,9 @@ describe('updateRow — partial merge', () => { }) it('allows explicitly nulling a field while preserving others', async () => { + dbChainMockFns.returning.mockResolvedValueOnce([ + { ...EXISTING_ROW, data: { name: 'Alice', age: null }, updatedAt: PERSISTED_UPDATED_AT }, + ]) const result = await updateRow( { tableId: 'tbl-1', rowId: 'row-1', data: { age: null }, workspaceId: 'ws-1' }, TABLE, @@ -199,6 +208,9 @@ describe('updateRow — partial merge', () => { }) it('handles a full-row update correctly (idempotent merge)', async () => { + dbChainMockFns.returning.mockResolvedValueOnce([ + { ...EXISTING_ROW, data: { name: 'Bob', age: 25 }, updatedAt: PERSISTED_UPDATED_AT }, + ]) const result = await updateRow( { tableId: 'tbl-1', rowId: 'row-1', data: { name: 'Bob', age: 25 }, workspaceId: 'ws-1' }, TABLE, diff --git a/apps/sim/lib/table/application/rows.test.ts b/apps/sim/lib/table/application/rows.test.ts index 9e9c9e0dc71..e8b7974d555 100644 --- a/apps/sim/lib/table/application/rows.test.ts +++ b/apps/sim/lib/table/application/rows.test.ts @@ -141,7 +141,14 @@ vi.mock('@/lib/table/rows/secret-provenance', () => ({ ), }), createUnknownTableRowSecretProvenance: () => ({ complete: false, columns: {} }), - loadTableRowSecretProvenance: mockLoadSecretProvenance, + TableRowProvenanceReader: class { + constructor(scope: unknown) { + mockLoadSecretProvenance(scope) + } + exportProvenance() { + return { version: 1, complete: true, entries: [] } + } + }, })) vi.mock('@/lib/table/validation', () => ({ @@ -676,7 +683,8 @@ describe('row query and upsert application semantics', () => { predicate: { field: 'column_name', op: 'eq', value: 'Ada' }, sort: { column_name: 'asc' }, }), - expect.any(String) + expect.any(String), + undefined ) }) @@ -821,7 +829,8 @@ describe('row query and upsert application semantics', () => { expect(mockQueryRows).toHaveBeenCalledWith( TABLE, expect.objectContaining({ offset: 100 }), - expect.any(String) + expect.any(String), + undefined ) }) @@ -833,14 +842,13 @@ describe('row query and upsert application semantics', () => { createdAt: new Date('2026-01-01'), updatedAt: new Date('2026-01-01'), } - const provenance = { complete: true, columns: {} } + const provenance = { version: 1, complete: true, entries: [] } mockQueryRows.mockResolvedValueOnce({ rows: [row], rowCount: 1, totalCount: null, nextCursor: null, }) - mockLoadSecretProvenance.mockResolvedValueOnce(provenance) const result = await queryTableRows.execute({ principal: PRINCIPAL, @@ -851,14 +859,14 @@ describe('row query and upsert application semantics', () => { }, }) - // Narrowed to the columns the row still holds: without `selectedValues` a - // stale sidecar entry for a dropped column rides along in the envelope, and - // the unmigrated `rows`/`query` routes have always narrowed here. - expect(mockLoadSecretProvenance).toHaveBeenCalledWith( - [{ id: row.id, updatedAt: row.updatedAt, selectedValues: row.data }], - { userId: 'user-1', workspaceId: TABLE.workspaceId } + expect(mockLoadSecretProvenance).toHaveBeenCalledWith({ + userId: 'user-1', + workspaceId: TABLE.workspaceId, + }) + expect(mockQueryRows.mock.calls[0][3]).toEqual( + expect.objectContaining({ exportProvenance: expect.any(Function) }) ) - expect(result.secretProvenance).toBe(provenance) + expect(result.secretProvenance).toEqual(provenance) }) it('audits only the authoritative deleted count and suppresses no-op audit', async () => { @@ -1488,7 +1496,8 @@ describe('opt-in per-cell run state', () => { expect(mockQueryRows).toHaveBeenCalledWith( TABLE, expect.objectContaining({ withExecutions: true }), - expect.any(String) + expect.any(String), + undefined ) }) diff --git a/apps/sim/lib/table/application/rows.ts b/apps/sim/lib/table/application/rows.ts index 4bce8ce9ad5..284e8a794ca 100644 --- a/apps/sim/lib/table/application/rows.ts +++ b/apps/sim/lib/table/application/rows.ts @@ -84,7 +84,7 @@ import { createExactEmptyTableRowSecretProvenance, createTableRowSecretProvenanceFromRegistry, createUnknownTableRowSecretProvenance, - loadTableRowSecretProvenance, + TableRowProvenanceReader, } from '@/lib/table/rows/secret-provenance' import type { FindRowMatch, RowWriteOptions } from '@/lib/table/rows/service' import { replaceTableRowsWithTx } from '@/lib/table/rows/service' @@ -151,27 +151,16 @@ interface TableResult { table: TableDefinition } -type TableRowsProvenance = Awaited> +type TableRowsProvenance = ReturnType -async function loadAuthorizedRowsProvenance( +function createAuthorizedRowsProvenanceReader( workspaceId: string, attributedUserId: string, - // The loader reads only id, updatedAt and the selected values, so a row - // without its executions sidecar is enough — see `TABLE_ROW_SIDECAR_SELECTION`. - rows: TableRowSummary[], include: boolean | undefined -): Promise { - if (!include) return undefined - return loadTableRowSecretProvenance( - // `selectedValues` narrows the sidecar to the columns the row still holds. - // Without it a stale entry for a dropped column rides along in the envelope, - // which is how the unmigrated `rows`/`query` routes have always behaved. - rows.map((row) => ({ id: row.id, updatedAt: row.updatedAt, selectedValues: row.data })), - { - userId: attributedUserId, - workspaceId, - } - ) +): TableRowProvenanceReader | undefined { + return include + ? new TableRowProvenanceReader({ userId: attributedUserId, workspaceId }) + : undefined } function requestId(input: TableScopedInput): string { @@ -492,6 +481,11 @@ export const queryTableRows = defineAuthorizedTableUseCase({ operation: tableOperations.queryRows, resolveContext: ({ input }: { input: QueryTableRowsInput }) => resolveActiveTableContext(input), async execute({ principal, input, context }): Promise { + const readProvenance = createAuthorizedRowsProvenanceReader( + context.workspaceId, + actorUserId(principal, context.billedAccountUserId), + input.includePersistedSecretProvenance + ) try { if (input.requireV2Feature) { const orgId = await getWorkspaceOrganizationId(context.workspaceId) @@ -580,17 +574,13 @@ export const queryTableRows = defineAuthorizedTableUseCase({ runStateBudgetBytes: TABLE_LIMITS.MAX_ROW_RUN_STATE_BYTES, columnIds, }, - requestId(input) + requestId(input), + readProvenance ) return { table: context.table, ...result, - secretProvenance: await loadAuthorizedRowsProvenance( - context.workspaceId, - actorUserId(principal, context.billedAccountUserId), - result.rows, - input.includePersistedSecretProvenance - ), + secretProvenance: readProvenance?.exportProvenance(), } } catch (error) { rethrowQueryValidation(error) @@ -655,7 +645,17 @@ export const readTableRow = defineAuthorizedTableUseCase({ operation: tableOperations.readRow, resolveContext: ({ input }: { input: ReadTableRowInput }) => resolveActiveTableContext(input), async execute({ principal, input, context }): Promise { - const row = await getRowSummaryById(context.tableId, input.rowId, context.workspaceId) + const readProvenance = createAuthorizedRowsProvenanceReader( + context.workspaceId, + actorUserId(principal, context.billedAccountUserId), + input.includePersistedSecretProvenance + ) + const row = await getRowSummaryById( + context.tableId, + input.rowId, + context.workspaceId, + readProvenance + ) if (!row) throw new OrchestrationError('not_found', 'Row not found') const runState = input.includeRunState ? await loadExecutionsForRow(db, input.rowId, { @@ -666,12 +666,7 @@ export const readTableRow = defineAuthorizedTableUseCase({ table: context.table, row, ...(runState ? { runState } : {}), - secretProvenance: await loadAuthorizedRowsProvenance( - context.workspaceId, - actorUserId(principal, context.billedAccountUserId), - [row], - input.includePersistedSecretProvenance - ), + secretProvenance: readProvenance?.exportProvenance(), } }, }) @@ -755,6 +750,11 @@ export const createTableRows = defineAuthorizedTableUseCase({ operation: tableOperations.createRows, resolveContext: ({ input }: { input: CreateTableRowsInput }) => resolveActiveTableContext(input), async execute({ principal, input, context }): Promise { + const readProvenance = createAuthorizedRowsProvenanceReader( + context.workspaceId, + actorUserId(principal, context.billedAccountUserId), + input.includePersistedSecretProvenance + ) const userId = actorUserId(principal, context.billedAccountUserId) if (input.kind === 'single') { if (input.afterRowId && input.beforeRowId) { @@ -774,7 +774,7 @@ export const createTableRows = defineAuthorizedTableUseCase({ input, storageData: data, }) - const writeOptions = rowWriteOptions(input) + const writeOptions = { ...rowWriteOptions(input), readProvenance } await throwValidationResponse( await validateRowData({ rowData: data, @@ -803,12 +803,7 @@ export const createTableRows = defineAuthorizedTableUseCase({ kind: 'single', table: context.table, row, - secretProvenance: await loadAuthorizedRowsProvenance( - context.workspaceId, - actorUserId(principal, context.billedAccountUserId), - [row], - input.includePersistedSecretProvenance - ), + secretProvenance: readProvenance?.exportProvenance(), } } if (input.rows.length < 1 || input.rows.length > TABLE_LIMITS.MAX_BATCH_INSERT_SIZE) { @@ -834,7 +829,7 @@ export const createTableRows = defineAuthorizedTableUseCase({ storageRows: rows, }).stamps : defaultedRowsSecretProvenance(rows, input.secretProvenance) - const batchWriteOptions = rowWriteOptions(input) + const batchWriteOptions = { ...rowWriteOptions(input), readProvenance } await throwValidationResponse( await validateBatchRows({ rows, @@ -861,12 +856,7 @@ export const createTableRows = defineAuthorizedTableUseCase({ kind: 'batch', table: context.table, rows: created, - secretProvenance: await loadAuthorizedRowsProvenance( - context.workspaceId, - actorUserId(principal, context.billedAccountUserId), - created, - input.includePersistedSecretProvenance - ), + secretProvenance: readProvenance?.exportProvenance(), } }, afterSuccess: ({ context, input, result }) => { @@ -1139,6 +1129,11 @@ export const updateTableRow = defineAuthorizedTableUseCase({ operation: tableOperations.updateRow, resolveContext: ({ input }: { input: UpdateTableRowInput }) => resolveActiveTableContext(input), async execute({ principal, input, context }): Promise { + const readProvenance = createAuthorizedRowsProvenanceReader( + context.workspaceId, + actorUserId(principal, context.billedAccountUserId), + input.includePersistedSecretProvenance + ) const data = rowDataToStorage(input.data, context.table, input.dataKeying, input.strictWrite) const secretProvenance = singleRowWriteProvenance({ principal, @@ -1159,19 +1154,14 @@ export const updateTableRow = defineAuthorizedTableUseCase({ }, context.table, requestId(input), - rowWriteOptions(input) + { ...rowWriteOptions(input), readProvenance } ) if (!row) throw new Error('Unconditional table row update was rejected') return { table: context.table, row, changed: Object.keys(data).length > 0, - secretProvenance: await loadAuthorizedRowsProvenance( - context.workspaceId, - actorUserId(principal, context.billedAccountUserId), - [row], - input.includePersistedSecretProvenance - ), + secretProvenance: readProvenance?.exportProvenance(), } }, afterSuccess: ({ context, input, result }) => { @@ -1457,6 +1447,11 @@ export const upsertTableRow = defineAuthorizedTableUseCase({ operation: tableOperations.upsertRow, resolveContext: ({ input }: { input: UpsertTableRowInput }) => resolveActiveTableContext(input), async execute({ principal, input, context }): Promise { + const readProvenance = createAuthorizedRowsProvenanceReader( + context.workspaceId, + actorUserId(principal, context.billedAccountUserId), + input.includePersistedSecretProvenance + ) // An id-keyed caller already names the storage column; only a name-keyed // one needs the lookup, and its miss falls through as before. const conflictTarget = @@ -1483,18 +1478,13 @@ export const upsertTableRow = defineAuthorizedTableUseCase({ }, context.table, requestId(input), - rowWriteOptions(input) + { ...rowWriteOptions(input), readProvenance } ) return { table: context.table, row: result.row, operation: result.operation, - secretProvenance: await loadAuthorizedRowsProvenance( - context.workspaceId, - actorUserId(principal, context.billedAccountUserId), - [result.row], - input.includePersistedSecretProvenance - ), + secretProvenance: readProvenance?.exportProvenance(), } }, afterSuccess: ({ context }) => signalTableRowsChanged(context.tableId), diff --git a/apps/sim/lib/table/planner.ts b/apps/sim/lib/table/planner.ts index 2220d5faa31..3461c712c16 100644 --- a/apps/sim/lib/table/planner.ts +++ b/apps/sim/lib/table/planner.ts @@ -53,12 +53,17 @@ async function setReadGuards(trx: DbTransaction, seqscanOff: boolean): Promise( fn: (trx: DbTransaction) => Promise, - opts?: { seqscanOff?: boolean } + opts?: { seqscanOff?: boolean; repeatableRead?: boolean } ): Promise { - return db.transaction(async (trx) => { - await setReadGuards(trx, opts?.seqscanOff ?? false) - return fn(trx) - }) + return db.transaction( + async (trx) => { + await setReadGuards(trx, opts?.seqscanOff ?? false) + return fn(trx) + }, + opts?.repeatableRead + ? { isolationLevel: 'repeatable read', accessMode: 'read only' } + : undefined + ) } /** diff --git a/apps/sim/lib/table/rows/ordering.ts b/apps/sim/lib/table/rows/ordering.ts index f20ebee0ef0..feca6ef446c 100644 --- a/apps/sim/lib/table/rows/ordering.ts +++ b/apps/sim/lib/table/rows/ordering.ts @@ -16,7 +16,10 @@ import type { MutationProof } from '@/lib/table/mutation-locks' import { keyBetween, nKeysBetween } from '@/lib/table/order-key' import { type DbExecutor, type DbTransaction, withSeqscanOff } from '@/lib/table/planner' import { TableRowNotFoundError } from '@/lib/table/rows/errors' -import { mutateTableRowsWithSecretProvenance } from '@/lib/table/rows/secret-provenance' +import { + mutateTableRowsWithSecretProvenance, + type TableRowProvenanceReader, +} from '@/lib/table/rows/secret-provenance' import { setTableTxTimeouts } from '@/lib/table/tx' import type { RowData, TableDefinition, TableRowSecretProvenanceWrite } from '@/lib/table/types' @@ -291,6 +294,7 @@ export async function resolveBatchInsertOrderKeys( * by the `increment_user_table_row_count` trigger. */ export async function insertOrderedRow(params: { + readProvenance?: TableRowProvenanceReader tableId: string workspaceId: string data: RowData @@ -337,7 +341,7 @@ export async function insertOrderedRow(params: { // order_key is authoritative — keep a best-effort, no-shift position. const targetPosition = await nextRowPosition(trx, tableId) - return mutateTableRowsWithSecretProvenance(trx, { + const rows = await mutateTableRowsWithSecretProvenance(trx, { rows: [{ rowId, provenance: secretProvenance }], rowState: 'new', mode: 'replace', @@ -362,6 +366,8 @@ export async function insertOrderedRow(params: { } }, }) + await params.readProvenance?.capture(trx, rows) + return rows }) return { id: row.id, diff --git a/apps/sim/lib/table/rows/secret-provenance.postgres.test.ts b/apps/sim/lib/table/rows/secret-provenance.postgres.test.ts index e6d60e3aa6b..d0a43d6558e 100644 --- a/apps/sim/lib/table/rows/secret-provenance.postgres.test.ts +++ b/apps/sim/lib/table/rows/secret-provenance.postgres.test.ts @@ -1,7 +1,7 @@ /** * @vitest-environment node * - * Runs the real Drizzle queries against temporary PostgreSQL tables. Set + * Runs the real Drizzle queries against an isolated PostgreSQL schema. Set * TABLE_PROVENANCE_TEST_DATABASE_URL to a local test database to include this suite. * From apps/sim, run: * `TABLE_PROVENANCE_TEST_DATABASE_URL=postgresql://user@127.0.0.1:5432/postgres bun run test lib/table/rows/secret-provenance.postgres.test.ts` @@ -10,20 +10,29 @@ */ import { userTableRows } from '@sim/db/schema' import { loggingSessionMock } from '@sim/testing' +import { generateShortId } from '@sim/utils/id' import { eq, sql } from 'drizzle-orm' import { drizzle, type PostgresJsDatabase } from 'drizzle-orm/postgres-js' import postgres from 'postgres' import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' import { PROVENANCE_MAX_SERIALIZED_BYTES } from '@/lib/execution/provenance-limits' import { createWorkflowCellProgressWriter } from '@/lib/table/cell-write' -import type { DbTransaction } from '@/lib/table/planner' +import { assertRowInsert } from '@/lib/table/mutation-locks' +import type { DbExecutor, DbTransaction } from '@/lib/table/planner' +import { insertOrderedRow } from '@/lib/table/rows/ordering' import { getTableSnapshotModelMountSafety, loadTableRowSecretProvenance, mutateTableRowsWithSecretProvenance, + TableRowProvenanceReader, updateTableRowsWithDerivedSecretProvenance, } from '@/lib/table/rows/secret-provenance' -import { getRowById, updateRow } from '@/lib/table/rows/service' +import { + fetchRowsBounded, + getRowById, + getRowSummaryById, + updateRow, +} from '@/lib/table/rows/service' import { fireTableTrigger } from '@/lib/table/trigger' import type { RowData, TableDefinition, TableRowSecretProvenanceWrite } from '@/lib/table/types' import { cancelWorkflowGroupRuns } from '@/lib/table/workflow-columns' @@ -88,7 +97,12 @@ const databaseUrl = process.env.TABLE_PROVENANCE_TEST_DATABASE_URL if (databaseUrl && !['localhost', '127.0.0.1', '[::1]'].includes(new URL(databaseUrl).hostname)) { throw new Error('Table provenance PostgreSQL tests require a local database') } -const connection = databaseUrl ? postgres(databaseUrl, { max: 1 }) : undefined +const testSchema = `provenance_${generateShortId() + .replace(/[^a-zA-Z0-9]/g, '') + .toLowerCase()}` +const connectionOptions = { max: 1, connection: { search_path: testSchema } } +const connection = databaseUrl ? postgres(databaseUrl, connectionOptions) : undefined +const writer = databaseUrl ? postgres(databaseUrl, connectionOptions) : undefined const updatedAt = new Date('2026-08-05T00:00:00.123Z') const secretEntry = { columnId: 'retained', encryptedValue: 'encrypted-secret', name: 'SECRET' } const scope = { userId: 'user-1', workspaceId: 'workspace-1' } @@ -186,16 +200,17 @@ async function writeWideRow() { describe.skipIf(!databaseUrl)('table provenance in PostgreSQL', () => { beforeAll(async () => { if (!connection) throw new Error('PostgreSQL test database is not initialized') + await connection`CREATE SCHEMA ${connection(testSchema)}` database.current = drizzle(connection) await connection.unsafe(` - CREATE TEMP TABLE user_table_definitions (id text PRIMARY KEY, workspace_id text NOT NULL, rows_version integer NOT NULL); - CREATE TEMP TABLE user_table_rows ( + CREATE TABLE user_table_definitions (id text PRIMARY KEY, workspace_id text NOT NULL, rows_version integer NOT NULL); + CREATE TABLE user_table_rows ( id text PRIMARY KEY, table_id text NOT NULL, workspace_id text NOT NULL, data jsonb NOT NULL, updated_at timestamp NOT NULL, secret_provenance_version integer, position integer NOT NULL DEFAULT 0, order_key text, created_at timestamp NOT NULL DEFAULT now(), created_by text ); - CREATE TEMP TABLE table_row_executions ( + CREATE TABLE table_row_executions ( table_id text NOT NULL, row_id text NOT NULL REFERENCES user_table_rows(id) ON DELETE CASCADE, group_id text NOT NULL, status text NOT NULL, execution_id text, job_id text, workflow_id text NOT NULL, error text, running_block_ids text[] NOT NULL DEFAULT '{}', @@ -203,11 +218,11 @@ describe.skipIf(!databaseUrl)('table provenance in PostgreSQL', () => { capability_governed_user_id text, enrichment_details jsonb, updated_at timestamp NOT NULL DEFAULT now(), PRIMARY KEY (row_id, group_id) ); - CREATE TEMP TABLE user_table_row_secret_provenance ( + CREATE TABLE user_table_row_secret_provenance ( row_id text PRIMARY KEY, content_updated_at timestamp NOT NULL, status text NOT NULL, entries jsonb NOT NULL, updated_at timestamp DEFAULT now() ); - CREATE FUNCTION pg_temp.demote_changed_row() RETURNS trigger LANGUAGE plpgsql AS $body$ + CREATE FUNCTION demote_changed_row() RETURNS trigger LANGUAGE plpgsql AS $body$ BEGIN IF NEW.data IS DISTINCT FROM OLD.data THEN NEW.updated_at := clock_timestamp(); @@ -217,7 +232,7 @@ describe.skipIf(!databaseUrl)('table provenance in PostgreSQL', () => { END $body$; CREATE TRIGGER demote_changed_row BEFORE UPDATE ON user_table_rows - FOR EACH ROW EXECUTE FUNCTION pg_temp.demote_changed_row(); + FOR EACH ROW EXECUTE FUNCTION demote_changed_row(); `) }) @@ -232,9 +247,140 @@ describe.skipIf(!databaseUrl)('table provenance in PostgreSQL', () => { }) afterAll(async () => { + await writer?.end() + if (connection) await connection`DROP SCHEMA ${connection(testSchema)} CASCADE` await connection?.end() }) + async function concurrentPatch(data: RowData) { + if (!writer) throw new Error('Writer unavailable') + await drizzle(writer).transaction(async (tx) => { + await mutateTableRowsWithSecretProvenance(tx as DbTransaction, { + rows: [{ rowId: 'row-1', provenance: { complete: true, columns: {} } }], + rowState: 'existing', + mode: 'replace', + mutate: async () => { + const rows = await tx + .update(userTableRows) + .set({ data }) + .where(eq(userTableRows.id, 'row-1')) + .returning() + return { value: rows, affectedRowIds: rows.map((row) => row.id) } + }, + }) + }) + } + + function interleaveWriter(reader: TableRowProvenanceReader) { + const capture = reader.capture.bind(reader) + vi.spyOn(reader, 'capture').mockImplementationOnce(async (executor, rows) => { + await concurrentPatch({ retained: 'new public value', removed: 'concurrent value' }) + await capture(executor, rows) + }) + } + + it('keeps the original single-row value and secret lineage across a concurrent stamped write', async () => { + mockIsEnforced.mockReturnValue(true) + await insertRow({ + status: 'exact', + entries: [ + { ...secretEntry, sourceUserId: scope.userId, sourceWorkspaceId: scope.workspaceId }, + ], + data: { retained: 'old secret value' }, + }) + const reader = new TableRowProvenanceReader(scope) + interleaveWriter(reader) + const row = await getRowSummaryById('table-1', 'row-1', 'workspace-1', reader) + expect(row?.data).toEqual({ retained: 'old secret value' }) + expect(reader.exportProvenance()).toEqual({ + version: 1, + complete: true, + scope, + entries: [{ encryptedValue: 'encrypted-secret', name: 'SECRET' }], + }) + expect((await getRowSummaryById('table-1', 'row-1', 'workspace-1'))?.data.retained).toBe( + 'new public value' + ) + }) + + it('captures only projected, returned rows and ignores an unknown pagination witness', async () => { + mockIsEnforced.mockReturnValue(true) + await insertRow({ status: 'exact', entries: [secretEntry] }) + await insertRow({ id: 'row-2', status: 'unknown' }) + const reader = new TableRowProvenanceReader(scope) + interleaveWriter(reader) + const result = await fetchRowsBounded({ + baseWhere: eq(userTableRows.tableId, 'table-1'), + orderBy: sql`${userTableRows.id} ASC`, + sorted: false, + keysetValid: false, + startOffset: 0, + limit: 1, + budgetBytes: 1024, + pageCutBytes: 1024, + columnIds: new Set(['removed']), + readProvenance: reader, + }) + expect(result.rows.map((row) => row.data)).toEqual([{ removed: 'other' }]) + expect(result.hasMore).toBe(true) + expect(reader.exportProvenance()).toEqual({ version: 1, complete: true, scope, entries: [] }) + }) + + it('captures insert provenance before a later writer changes the returned row', async () => { + mockIsEnforced.mockReturnValue(true) + const reader = new TableRowProvenanceReader(scope) + const row = await insertOrderedRow({ + tableId: 'table-1', + workspaceId: 'workspace-1', + rowId: 'row-1', + data: { retained: 'inserted' }, + now: updatedAt, + proof: assertRowInsert(table), + secretProvenance: { complete: true, columns: {} }, + readProvenance: reader, + }) + expect(reader.exportProvenance()).toEqual({ version: 1, complete: true, scope, entries: [] }) + await concurrentPatch({ retained: 'later' }) + expect(row.data.retained).toBe('inserted') + expect(reader.exportProvenance().complete).toBe(true) + }) + + it('returns the actual merged row and captures its lineage before the write transaction commits', async () => { + mockIsEnforced.mockReturnValue(true) + await insertRow({ status: 'exact' }) + if (!database.current) throw new Error('Database unavailable') + const transact = database.current.transaction.bind(database.current) + vi.spyOn(database.current, 'transaction').mockImplementationOnce(async (callback, config) => { + await concurrentPatch({ retained: 'value', removed: 'new untouched value' }) + return transact(callback, config) + }) + const reader = new TableRowProvenanceReader(scope) + const capture = reader.capture.bind(reader) + vi.spyOn(reader, 'capture').mockImplementation(async (executor: DbExecutor, rows) => { + expect(executor).not.toBe(database.current) + const [outside] = await writer!`SELECT data FROM user_table_rows WHERE id = 'row-1'` + expect(outside.data.retained).toBe('value') + await capture(executor, rows) + }) + const result = await updateRow( + { + tableId: 'table-1', + workspaceId: 'workspace-1', + rowId: 'row-1', + data: { retained: 'patched' }, + secretProvenance: { complete: true, columns: {} }, + }, + table, + 'request-1', + { readProvenance: reader } + ) + expect(result?.data).toEqual({ retained: 'patched', removed: 'new untouched value' }) + expect(reader.exportProvenance()).toEqual({ version: 1, complete: true, scope, entries: [] }) + await concurrentPatch({ retained: 'later' }) + expect(result?.data.retained).toBe('patched') + expect(reader.exportProvenance().complete).toBe(true) + }) + it.each(['exact', 'unknown', 'legacy', 'stale'] as const)( 'preserves %s row provenance through cancellation, restart, and a cell write', async (baseStatus) => { diff --git a/apps/sim/lib/table/rows/secret-provenance.test.ts b/apps/sim/lib/table/rows/secret-provenance.test.ts index 1cbcd3b7a48..36dc1b8f1dd 100644 --- a/apps/sim/lib/table/rows/secret-provenance.test.ts +++ b/apps/sim/lib/table/rows/secret-provenance.test.ts @@ -29,6 +29,7 @@ import { getTableSnapshotModelMountSafety, loadTableRowSecretProvenance, mutateTableRowsWithSecretProvenance, + TableRowProvenanceReader, updateTableRowsWithDerivedSecretProvenance, } from '@/lib/table/rows/secret-provenance' @@ -213,6 +214,75 @@ describe('table row secret provenance', () => { }) }) + it('reports a row revision mismatch without exposing row content', async () => { + queueTableRows(userTableRows, [ + { + id: 'tracked-row', + updatedAt: new Date(ROW_UPDATED_AT.getTime() + 1), + secretProvenanceVersion: 1, + sidecarStatus: 'exact', + sidecarEntries: [], + sidecarIsCurrent: true, + }, + ]) + const provenance = await loadTableRowSecretProvenance( + [{ id: 'tracked-row', updatedAt: ROW_UPDATED_AT }], + { + userId: 'user-1', + workspaceId: 'workspace-1', + } + ) + expect(provenance.complete).toBe(false) + expect(mockError).toHaveBeenCalledWith('Table row read could not establish secret provenance', { + surface: 'table-row', + cause: 'row-revision-mismatch', + rowCount: 1, + workspaceId: 'workspace-1', + actorUserId: 'user-1', + }) + }) + + it.each([{ selectedColumns: ['input-column'] }, { selectedColumns: [] }])( + 'captures only selected worker input columns: %j', + async ({ selectedColumns }) => { + queueTableRows(userTableRows, [ + { + id: 'tracked-row', + updatedAt: ROW_UPDATED_AT, + secretProvenanceVersion: 1, + sidecarStatus: 'exact', + sidecarIsCurrent: true, + sidecarEntries: ['input-column', 'output-column'].map((columnId) => ({ + columnId, + encryptedValue: `encrypted-${columnId}`, + name: columnId, + sourceUserId: 'user-1', + sourceWorkspaceId: 'workspace-1', + })), + }, + ]) + const reader = new TableRowProvenanceReader( + { userId: 'user-1', workspaceId: 'workspace-1' }, + new Set(selectedColumns) + ) + + await reader.capture(dbChainMock.db as unknown as DbTransaction, [ + { + id: 'tracked-row', + updatedAt: ROW_UPDATED_AT, + data: { 'input-column': 'input-secret', 'output-column': 'output-secret' }, + }, + ]) + + expect(reader.exportProvenance()).toEqual({ + version: 1, + complete: true, + scope: { userId: 'user-1', workspaceId: 'workspace-1' }, + entries: selectedColumns.map((name) => ({ name, encryptedValue: `encrypted-${name}` })), + }) + } + ) + it('does not activate a secret from an unselected column with the same value', async () => { queueTableRows(userTableRows, [ { diff --git a/apps/sim/lib/table/rows/secret-provenance.ts b/apps/sim/lib/table/rows/secret-provenance.ts index 604ef0c8b5e..dbd50ee65bc 100644 --- a/apps/sim/lib/table/rows/secret-provenance.ts +++ b/apps/sim/lib/table/rows/secret-provenance.ts @@ -20,6 +20,7 @@ import type { DbExecutor, DbTransaction } from '@/lib/table/planner' import type { RowData, TableRowSecretProvenanceWrite } from '@/lib/table/types' import { isResolvedSecretTraceProvenanceV1, + ResolvedSecretTraceProvenanceAccumulator, type ResolvedSecretTraceProvenanceEntryV1, type ResolvedSecretTraceProvenanceV1, ResolvedSecretTraceRegistry, @@ -896,17 +897,29 @@ function createStoredEntryAggregator(scope: ResolvedSecretTraceScopeV1) { */ export async function loadTableRowSecretProvenance( rows: TableRowCrossing[], - scope: ResolvedSecretTraceScopeV1 + scope: ResolvedSecretTraceScopeV1, + executor: DbExecutor = db ): Promise { if (rows.length === 0) { return { version: 1, complete: true, entries: [], scope } } + const incomplete = (cause: string): ResolvedSecretTraceProvenanceV1 => { + logger.error('Table row read could not establish secret provenance', { + surface: 'table-row', + cause, + rowCount: rows.length, + workspaceId: scope.workspaceId, + actorUserId: scope.userId, + }) + return { version: 1, complete: false, entries: [], scope } + } + const crossingById = new Map }>() for (const row of rows) { const existing = crossingById.get(row.id) if (existing && !sameTimestamp(existing.updatedAt, row.updatedAt)) { - return { version: 1, complete: false, entries: [], scope } + return incomplete('duplicate-row-revision') } const selectedColumnIds = row.selectedValues ? new Set(Object.keys(row.selectedValues)) @@ -922,7 +935,7 @@ export async function loadTableRowSecretProvenance( for (const columnId of selectedColumnIds) existing.selectedColumnIds.add(columnId) } const rowIds = [...crossingById.keys()] - const currentRows = await selectRowsWithSidecars(db, rowIds) + const currentRows = await selectRowsWithSidecars(executor, rowIds) const currentById = new Map(currentRows.map((row) => [row.id, row])) const aggregator = createStoredEntryAggregator(scope) @@ -930,8 +943,9 @@ export async function loadTableRowSecretProvenance( for (const rowId of rowIds) { const current = currentById.get(rowId) const crossing = crossingById.get(rowId) - if (!current || !crossing || !sameTimestamp(current.updatedAt, crossing.updatedAt)) { - return { version: 1, complete: false, entries: [], scope } + if (!current || !crossing) return incomplete('row-missing') + if (!sameTimestamp(current.updatedAt, crossing.updatedAt)) { + return incomplete('row-revision-mismatch') } if (current.secretProvenanceVersion === null) continue if ( @@ -945,16 +959,16 @@ export async function loadTableRowSecretProvenance( * read the table. Unenforced, the row contributes nothing, exactly like the legacy row above. */ if (isDurableSecretProvenanceEnforced('table-row')) { - return { version: 1, complete: false, entries: [], scope } + return incomplete('row-sidecar-not-exact') } unrecordedRowCount += 1 continue } const parsed = normalizeStoredEntries(current.sidecarEntries) - if (!parsed) return { version: 1, complete: false, entries: [], scope } + if (!parsed) return incomplete('row-sidecar-malformed') for (const entry of parsed) { if (crossing.selectedColumnIds && !crossing.selectedColumnIds.has(entry.columnId)) continue - if (!aggregator.add(entry)) return { version: 1, complete: false, entries: [], scope } + if (!aggregator.add(entry)) return incomplete('row-provenance-budget-exceeded') } } @@ -969,7 +983,7 @@ export async function loadTableRowSecretProvenance( } const entries = aggregator.build() - if (!entries) return { version: 1, complete: false, entries: [], scope } + if (!entries) return incomplete('row-provenance-budget-exceeded') const provenance: ResolvedSecretTraceProvenanceV1 = { version: 1, complete: true, @@ -978,5 +992,48 @@ export async function loadTableRowSecretProvenance( } return isResolvedSecretTraceProvenanceV1(provenance) ? provenance - : { version: 1, complete: false, entries: [], scope } + : incomplete('row-provenance-budget-exceeded') +} + +/** + * Collects only returned row values while their database snapshot is still valid. + * Readers use one repeatable-read transaction per bounded batch; writers capture + * after stamping and before releasing row locks. Nothing is reloaded after commit. + */ +export class TableRowProvenanceReader { + private readonly accumulator: ResolvedSecretTraceProvenanceAccumulator + + constructor( + private readonly scope: ResolvedSecretTraceScopeV1, + private readonly selectedColumnIds?: ReadonlySet + ) { + this.accumulator = new ResolvedSecretTraceProvenanceAccumulator(scope) + } + + async capture( + executor: DbExecutor, + rows: { id: string; updatedAt: Date | string; data: unknown }[] + ): Promise { + this.accumulator.record( + await loadTableRowSecretProvenance( + rows.map((row) => { + const data = row.data as RowData + let selectedValues = data + if (this.selectedColumnIds) { + selectedValues = {} + for (const columnId of this.selectedColumnIds) { + if (Object.hasOwn(data, columnId)) selectedValues[columnId] = data[columnId] + } + } + return { id: row.id, updatedAt: row.updatedAt, selectedValues } + }), + this.scope, + executor + ) + ) + } + + exportProvenance(): ResolvedSecretTraceProvenanceV1 { + return this.accumulator.exportProvenance() + } } diff --git a/apps/sim/lib/table/rows/service.ts b/apps/sim/lib/table/rows/service.ts index 29075ef47bb..eb047b0fb2c 100644 --- a/apps/sim/lib/table/rows/service.ts +++ b/apps/sim/lib/table/rows/service.ts @@ -62,7 +62,10 @@ import { selectRowIdPage, } from '@/lib/table/rows/ordering' import { pendingDeleteMask } from '@/lib/table/rows/pending-delete-mask' -import { mutateTableRowsWithSecretProvenance } from '@/lib/table/rows/secret-provenance' +import { + mutateTableRowsWithSecretProvenance, + type TableRowProvenanceReader, +} from '@/lib/table/rows/secret-provenance' import { buildFilterClause, buildPredicateClause, @@ -139,6 +142,7 @@ async function dispatchDeleteTriggers( * @throws Error if validation fails or capacity exceeded */ export interface RowWriteOptions { + readProvenance?: TableRowProvenanceReader /** * What this write does with a value its column's type cannot coerce. Defaults * to `null` — the cell is blanked and the write succeeds, which is what every @@ -203,6 +207,7 @@ export async function insertRow( now, secretProvenance: data.secretProvenance, proof: insertProof, + readProvenance: options.readProvenance, }) notifyTableRowUsage({ @@ -389,6 +394,7 @@ export async function batchInsertRowsWithTx( updatedAt: r.updatedAt, })) + await options.readProvenance?.capture(trx, result) return result } @@ -823,6 +829,7 @@ export async function upsertRow( }, }) if (!updatedRow) throw new Error('Matched table row no longer exists') + await options.readProvenance?.capture(trx, [updatedRow]) // No executions sidecar: no upsert surface puts one on the wire, and // loading it here would hold the write transaction open for a result that @@ -871,6 +878,7 @@ export async function upsertRow( }, }) if (!insertedRow) throw new Error('Failed to insert table row') + await options.readProvenance?.capture(trx, [insertedRow]) return { row: { @@ -1136,7 +1144,8 @@ async function countRowsTenantBounded(whereClause: SQL | undefined): Promise { const { filter, @@ -1224,6 +1233,7 @@ export async function queryRows( budgetBytes: TABLE_LIMITS.MAX_QUERY_RESULT_BYTES, pageCutBytes: getMaxPageBytes(), columnIds, + readProvenance, }) const [fetched, totalCount] = await Promise.all([drainPromise, countPromise]) @@ -1291,6 +1301,7 @@ export async function queryRows( } export interface BoundedFetchParams { + readProvenance?: TableRowProvenanceReader /** Tenant + delete-mask + user filter — WITHOUT any seek predicate. */ baseWhere: SQL | undefined orderBy: SQL @@ -1430,70 +1441,71 @@ export async function fetchRowsBounded(params: BoundedFetchParams): Promise 0 ? query.offset(batchOffset) : query } - // One tx per batch (SET LOCAL dies with it; holding a tx across JS - // accounting between batches would pin a pooled connection). Custom sorts - // order by `data->>'col'` — unestimatable — so they also penalize seq scans - // (9.7s→0.76s on a 1M-row table); default-order pages stream the index and - // just need the read timeout. Either way the batch runs under a statement - // timeout so a pathological filter can't scan unbounded. - return withReadGuards(async (trx) => buildQuery(trx), { seqscanOff: sorted }) + return withReadGuards( + async (trx) => { + const batch = await buildQuery(trx) + let cut = false + const returnedRows: Array = [] + for (const fetchedRow of batch) { + // Project before measuring: the budget is a promise about the response, + // so columns the caller will never receive must not count against it. + const row = columnIds + ? { ...fetchedRow, data: projectRowData(fetchedRow.data as RowData, columnIds) } + : fetchedRow + const rowBytes = Buffer.byteLength(JSON.stringify(row.data)) + const rowStoredBytes = columnIds + ? Buffer.byteLength(JSON.stringify(fetchedRow.data)) + : rowBytes + if (cutBytes !== undefined && rows.length > 0 && bytes + rowBytes > cutBytes) { + // Unbounded queries promise the ENTIRE result — a partial page would be + // silent truncation, so fail fast instead (the drain has only fetched + // ~budget bytes at this point, never the whole table). + if (limit === undefined) { + throw new TableQueryValidationError( + `Query result exceeds the ${Math.floor(cutBytes / (1024 * 1024))}MB limit. Add a filter or a limit to narrow the result.`, + 'TABLE_QUERY_RESULT_TOO_LARGE' + ) + } + // Bounded page, byte cut opted in: `row` is the witness. Requires a + // non-empty page so a single over-budget row is still returned alone. + hasMore = true + cut = true + break + } + // Limit cut: `row` is the +1 peek witness. + if (rows.length === limit) { + hasMore = true + cut = true + break + } + rows.push(row) + returnedRows.push(row) + bytes += rowBytes + storedBytes += rowStoredBytes + consumedSinceAnchor++ + if (rowBytes > maxRowBytes) maxRowBytes = rowBytes + if (rowStoredBytes > maxStoredRowBytes) maxStoredRowBytes = rowStoredBytes + if (keysetValid && row.orderKey) { + anchor = { orderKey: row.orderKey, id: row.id } + anchorOffset = 0 + consumedSinceAnchor = 0 + } + } + await params.readProvenance?.capture(trx, returnedRows) + return { cut, batchLength: batch.length } + }, + { seqscanOff: sorted, repeatableRead: Boolean(params.readProvenance) } + ) } while (true) { const limitRemaining = limit === undefined ? Number.POSITIVE_INFINITY : limit - rows.length const target = Math.min(nextBatchRows(), limitRemaining) - const ask = target + 1 // +1 = witness row proving more data exists past a cut - const batch = await runBatch(anchor, anchorOffset + consumedSinceAnchor, ask) - if (batch.length === 0) break - - let cut = false - for (const fetchedRow of batch) { - // Project before measuring: the budget is a promise about the response, - // so columns the caller will never receive must not count against it. - const row = columnIds - ? { ...fetchedRow, data: projectRowData(fetchedRow.data as RowData, columnIds) } - : fetchedRow - const rowBytes = Buffer.byteLength(JSON.stringify(row.data)) - const rowStoredBytes = columnIds - ? Buffer.byteLength(JSON.stringify(fetchedRow.data)) - : rowBytes - if (cutBytes !== undefined && rows.length > 0 && bytes + rowBytes > cutBytes) { - // Unbounded queries promise the ENTIRE result — a partial page would be - // silent truncation, so fail fast instead (the drain has only fetched - // ~budget bytes at this point, never the whole table). - if (limit === undefined) { - throw new TableQueryValidationError( - `Query result exceeds the ${Math.floor(cutBytes / (1024 * 1024))}MB limit. Add a filter or a limit to narrow the result.`, - 'TABLE_QUERY_RESULT_TOO_LARGE' - ) - } - // Bounded page, byte cut opted in: `row` is the witness. Requires a - // non-empty page so a single over-budget row is still returned alone. - hasMore = true - cut = true - break - } - // Limit cut: `row` is the +1 peek witness. - if (rows.length === limit) { - hasMore = true - cut = true - break - } - rows.push(row) - bytes += rowBytes - storedBytes += rowStoredBytes - consumedSinceAnchor++ - if (rowBytes > maxRowBytes) maxRowBytes = rowBytes - if (rowStoredBytes > maxStoredRowBytes) maxStoredRowBytes = rowStoredBytes - if (keysetValid && row.orderKey) { - anchor = { orderKey: row.orderKey, id: row.id } - anchorOffset = 0 - consumedSinceAnchor = 0 - } - } + const ask = target + 1 + const { cut, batchLength } = await runBatch(anchor, anchorOffset + consumedSinceAnchor, ask) if (cut) break // Short batch = the source is exhausted; hasMore stays false. - if (batch.length < ask) break + if (batchLength < ask) break } return { @@ -1508,8 +1520,13 @@ export async function fetchRowsBounded(params: BoundedFetchParams): Promise -function selectRowRecord(tableId: string, rowId: string, workspaceId: string) { - return db +function selectRowRecord( + tableId: string, + rowId: string, + workspaceId: string, + executor: DbExecutor = db +) { + return executor .select() .from(userTableRows) .where( @@ -1552,10 +1569,16 @@ function toRowSummary(row: Awaited>[number]): export async function getRowSummaryById( tableId: string, rowId: string, - workspaceId: string + workspaceId: string, + readProvenance?: TableRowProvenanceReader ): Promise { - const [row] = await selectRowRecord(tableId, rowId, workspaceId) - return row ? toRowSummary(row) : null + const read = async (executor: DbExecutor) => { + const [row] = await selectRowRecord(tableId, rowId, workspaceId, executor) + const result = row ? toRowSummary(row) : null + if (result) await readProvenance?.capture(executor, [result]) + return result + } + return readProvenance ? withReadGuards(read, { repeatableRead: true }) : read(db) } /** One row with its executions sidecar, for the write and background paths. */ @@ -1662,6 +1685,16 @@ export async function updateRow( throw new TableRowNotFoundError() } if (Object.keys(data.data).length === 0 && data.executionsPatch === undefined) { + if (options.readProvenance) { + const row = await getRowSummaryById( + data.tableId, + data.rowId, + data.workspaceId, + options.readProvenance + ) + if (!row) throw new TableRowNotFoundError() + return { ...row, executions: existingRow.executions } + } return existingRow } @@ -1748,16 +1781,15 @@ export async function updateRow( // commit in one transaction so a partial write can't leave the sidecar // and the row out of sync. const guard = data.cancellationGuard - let persistedUpdatedAt: Date + let persistedRow: typeof userTableRows.$inferSelect try { - persistedUpdatedAt = await db.transaction(async (trx) => { + persistedRow = await db.transaction(async (trx) => { const mutate = async () => { const condition = and( eq(userTableRows.id, data.rowId), eq(userTableRows.tableId, data.tableId), eq(userTableRows.workspaceId, data.workspaceId) ) - const projection = { id: userTableRows.id, updatedAt: userTableRows.updatedAt } /** * Execution metadata has its own sidecar clock. Lock the content row for * existence and cancellation atomicity without invalidating its provenance. @@ -1768,8 +1800,8 @@ export async function updateRow( .update(userTableRows) .set({ data: persistedData, updatedAt: now }) .where(condition) - .returning(projection) - : await trx.select(projection).from(userTableRows).where(condition).for('update') + .returning() + : await trx.select().from(userTableRows).where(condition).for('update') if (!updatedRow) throw new TableRowNotFoundError() const result = await writeExecutionsPatch( @@ -1783,19 +1815,22 @@ export async function updateRow( throw new GuardRejected() } return { - value: updatedRow.updatedAt, + value: updatedRow, affectedRowIds: [updatedRow.id], } } - if (patchedColumnIds.size === 0) return (await mutate()).value - - return await mutateTableRowsWithSecretProvenance(trx, { - rows: [{ rowId: data.rowId, provenance: data.secretProvenance }], - rowState: 'existing', - mode: 'merge', - mutate, - }) + const row = + patchedColumnIds.size === 0 + ? (await mutate()).value + : await mutateTableRowsWithSecretProvenance(trx, { + rows: [{ rowId: data.rowId, provenance: data.secretProvenance }], + rowState: 'existing', + mode: 'merge', + mutate, + }) + await options.readProvenance?.capture(trx, [row]) + return row }) } catch (err) { if (err instanceof GuardRejected) return null @@ -1805,12 +1840,8 @@ export async function updateRow( logger.info(`[${requestId}] Updated row ${data.rowId} in table ${data.tableId}`) const updatedRow: TableRow = { - id: data.rowId, - data: mergedData, + ...toRowSummary(persistedRow), executions: mergedExecutions, - position: existingRow.position, - createdAt: existingRow.createdAt, - updatedAt: persistedUpdatedAt, } if (patchedColumnIds.size === 0) return updatedRow diff --git a/apps/sim/lib/uploads/contexts/workspace/workspace-file-secret-provenance.ts b/apps/sim/lib/uploads/contexts/workspace/workspace-file-secret-provenance.ts index 88d1b89f889..d3afdaef602 100644 --- a/apps/sim/lib/uploads/contexts/workspace/workspace-file-secret-provenance.ts +++ b/apps/sim/lib/uploads/contexts/workspace/workspace-file-secret-provenance.ts @@ -739,6 +739,15 @@ export async function getBoundWorkspaceFileSecretProvenance( workspaceId: string, identity: WorkspaceFileSecretProvenanceIdentity ): Promise { + /** + * Saving a chat upload changes its context but preserves its id, key, and bytes. Only a reader + * that captured the content revision may follow that promotion; the revision check below still + * refuses an intervening content write, including on a legacy untracked file. + */ + const contextCondition = + identity.context === 'mothership' && identity.contentUpdatedAt + ? inArray(workspaceFiles.context, ['mothership', 'workspace']) + : eq(workspaceFiles.context, identity.context) const [row] = await db .select({ fileContentUpdatedAt: workspaceFiles.contentUpdatedAt, @@ -757,7 +766,7 @@ export async function getBoundWorkspaceFileSecretProvenance( eq(workspaceFiles.id, identity.fileId), eq(workspaceFiles.key, identity.key), eq(workspaceFiles.workspaceId, workspaceId), - eq(workspaceFiles.context, identity.context) + contextCondition // Deliberately no deletedAt filter: `id` alone pins the exact row, and // recently-deleted/ reads are a real surface — excluding soft-deleted // rows made every archived file read as provenance-unknown and refused. diff --git a/apps/sim/tools/index.test.ts b/apps/sim/tools/index.test.ts index 17dde72d10e..e2c0fb8508e 100644 --- a/apps/sim/tools/index.test.ts +++ b/apps/sim/tools/index.test.ts @@ -28,6 +28,7 @@ import { DrizzleQueryError } from 'drizzle-orm/errors' import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' import { projectToolResultForCopilot } from '@/lib/copilot/request/tools/resolved-secret-result' +import type { EnvironmentResolutionSnapshot } from '@/lib/environment/utils' import { executeBitbucketTool } from '@/lib/internal/bitbucket/execute-tool' import { createInternalToolFileResult } from '@/lib/internal/tool-operations/file-result' import type { InternalToolOperationCall } from '@/lib/internal/tool-operations/types' @@ -93,7 +94,8 @@ const { const mockSecureFetchWithPinnedIP = inputValidationMockFns.mockSecureFetchWithPinnedIP const mockValidateUrlWithDNS = inputValidationMockFns.mockValidateUrlWithDNS -const mockGetEffectiveDecryptedEnv = environmentUtilsMockFns.mockGetEffectiveDecryptedEnv +const mockGetEffectiveEnvironmentSnapshot = + environmentUtilsMockFns.mockGetEffectiveEnvironmentSnapshot // Mock getBYOKKey vi.mock('@/lib/api-key/byok', () => ({ @@ -5041,6 +5043,21 @@ describe('Managed OAuth Credential Delegation', () => { describe('Copilot Env Variable Reference Resolution', () => { let cleanupEnvVars: () => void + function environmentSnapshot(variables: Record): EnvironmentResolutionSnapshot { + return { + personalEncrypted: {}, + workspaceEncrypted: Object.fromEntries( + Object.keys(variables).map((name) => [name, `encrypted-${name}`]) + ), + personalDecrypted: {}, + workspaceDecrypted: variables, + personalOwners: {}, + conflicts: [], + decryptionFailures: [], + workspaceUnredactedKeys: [], + } + } + function sentOperationInput(): Record { const call = mockExecuteInternalToolOperation.mock.calls.findLast( ([request]) => request.toolId === 'test_env_ref_tool' @@ -5069,8 +5086,10 @@ describe('Copilot Env Variable Reference Resolution', () => { NEXT_PUBLIC_APP_URL: 'http://localhost:3000', INTERNAL_API_BASE_URL: '', }) - mockGetEffectiveDecryptedEnv.mockReset() - mockGetEffectiveDecryptedEnv.mockResolvedValue({ SENTRY_AUTH_TOKEN: 'sntrys_real_token' }) + mockGetEffectiveEnvironmentSnapshot.mockReset() + mockGetEffectiveEnvironmentSnapshot.mockResolvedValue( + environmentSnapshot({ SENTRY_AUTH_TOKEN: 'sntrys_real_token' }) + ) }) afterEach(() => { @@ -5086,32 +5105,34 @@ describe('Copilot Env Variable Reference Resolution', () => { ) expect(result.success).toBe(true) - expect(mockGetEffectiveDecryptedEnv).toHaveBeenCalledWith('user-123', 'workspace-456') + expect(mockGetEffectiveEnvironmentSnapshot).toHaveBeenCalledWith('user-123', 'workspace-456') expect(sentOperationInput().apiKey).toBe('sntrys_real_token') }) it('keeps direct integration execution raw while projecting only its active workspace secret', async () => { const activeSecret = 'xxxxxxxx' const unusedSecret = 'true' - mockGetEffectiveDecryptedEnv.mockResolvedValueOnce({ - SERPER_API_KEY: activeSecret, - UNUSED_SECRET: unusedSecret, - }) + mockGetEffectiveEnvironmentSnapshot.mockResolvedValueOnce( + environmentSnapshot({ SERPER_API_KEY: activeSecret, UNUSED_SECRET: unusedSecret }) + ) mockExecuteInternalToolOperation.mockResolvedValueOnce( Response.json({ reflected: activeSecret, ordinary: unusedSecret }) ) - const registry = new ResolvedSecretTraceRegistry([ - { - name: 'SERPER_API_KEY', - plaintext: activeSecret, - encryptedValue: 'encrypted-active', - }, - { - name: 'UNUSED_SECRET', - plaintext: unusedSecret, - encryptedValue: 'encrypted-unused', - }, - ]) + const registry = new ResolvedSecretTraceRegistry( + [ + { + name: 'SERPER_API_KEY', + plaintext: activeSecret, + encryptedValue: 'encrypted-active', + }, + { + name: 'UNUSED_SECRET', + plaintext: unusedSecret, + encryptedValue: 'encrypted-unused', + }, + ], + { userId: 'user-123', workspaceId: 'workspace-456' } + ) const callerParams = { apiKey: '{{SERPER_API_KEY}}' } const result = await executeTool('test_env_ref_tool', callerParams, { @@ -5124,7 +5145,7 @@ describe('Copilot Env Variable Reference Resolution', () => { output: { reflected: activeSecret, ordinary: unusedSecret }, }) expect(callerParams).toEqual({ apiKey: '{{SERPER_API_KEY}}' }) - expect(mockGetEffectiveDecryptedEnv).toHaveBeenCalledWith('user-123', 'workspace-456') + expect(mockGetEffectiveEnvironmentSnapshot).toHaveBeenCalledWith('user-123', 'workspace-456') expect(sentOperationInput().apiKey).toBe(activeSecret) expect(projectToolResultForCopilot(result, registry)).toMatchObject({ success: true, @@ -5132,23 +5153,64 @@ describe('Copilot Env Variable Reference Resolution', () => { }) }) + it('tracks a secret added after the Copilot turn started', async () => { + const secret = 'new-workspace-api-key' + environmentUtilsMockFns.mockGetEffectiveEnvironmentSnapshot.mockResolvedValueOnce({ + personalEncrypted: {}, + workspaceEncrypted: { ADDED_API_KEY: 'encrypted-new-key' }, + personalDecrypted: {}, + workspaceDecrypted: { ADDED_API_KEY: secret }, + personalOwners: {}, + conflicts: [], + decryptionFailures: [], + workspaceUnredactedKeys: [], + }) + const parent = new ResolvedSecretTraceRegistry([], { + userId: 'user-123', + workspaceId: 'workspace-456', + }) + const registry = parent.forkForInputPaths([]) + mockExecuteInternalToolOperation.mockResolvedValueOnce(Response.json({ reflected: secret })) + + const result = await executeTool( + 'test_env_ref_tool', + { apiKey: '{{ADDED_API_KEY}}' }, + { executionContext: copilotContext(), resolvedSecretTraceRegistry: registry } + ) + + expect(sentOperationInput().apiKey).toBe(secret) + expect(projectToolResultForCopilot(result, registry)).toMatchObject({ + success: true, + output: { reflected: '{{ADDED_API_KEY}}' }, + }) + expect(registry.isComplete()).toBe(true) + expect(parent.getActiveMatches()).toEqual([]) + parent.mergeToolCallRegistry(registry) + expect(parent.getActiveMatches()).toEqual([ + { plaintext: secret, replacement: '{{ADDED_API_KEY}}' }, + ]) + }) + it('does not let a pending user-only reference affect an unrelated result', async () => { const secret = 'sntrys_real_token' - const registry = new ResolvedSecretTraceRegistry([ - { - name: 'SENTRY_AUTH_TOKEN', - plaintext: secret, - encryptedValue: 'encrypted-token', - }, - ]) - let resolveEnvironment!: (variables: Record) => void + const registry = new ResolvedSecretTraceRegistry( + [ + { + name: 'SENTRY_AUTH_TOKEN', + plaintext: secret, + encryptedValue: 'encrypted-token', + }, + ], + { userId: 'user-123', workspaceId: 'workspace-456' } + ) + let resolveEnvironment!: (environment: EnvironmentResolutionSnapshot) => void let markResolutionStarted!: () => void const resolutionStarted = new Promise((resolve) => { markResolutionStarted = resolve }) - mockGetEffectiveDecryptedEnv.mockImplementationOnce( + mockGetEffectiveEnvironmentSnapshot.mockImplementationOnce( () => - new Promise>((resolve) => { + new Promise((resolve) => { resolveEnvironment = resolve markResolutionStarted() }) @@ -5168,7 +5230,7 @@ describe('Copilot Env Variable Reference Resolution', () => { projectToolResultForCopilot({ success: true, output: { result: secret } }, registry) ).toMatchObject({ output: { result: secret } }) - resolveEnvironment({ SENTRY_AUTH_TOKEN: secret }) + resolveEnvironment(environmentSnapshot({ SENTRY_AUTH_TOKEN: secret })) await expect(execution).resolves.toMatchObject({ success: true }) expect(registry.isComplete()).toBe(true) @@ -5208,7 +5270,7 @@ describe('Copilot Env Variable Reference Resolution', () => { expect(result.success).toBe(true) expect(sentOperationInput().apiKey).toBe('Bearer {{SENTRY_AUTH_TOKEN}}') - expect(mockGetEffectiveDecryptedEnv).not.toHaveBeenCalled() + expect(mockGetEffectiveEnvironmentSnapshot).not.toHaveBeenCalled() }) it('fails with a clear error before any request when the variable is missing', async () => { @@ -5239,7 +5301,7 @@ describe('Copilot Env Variable Reference Resolution', () => { expect(result.success).toBe(false) expect(result.error).toContain('authenticated user context') - expect(mockGetEffectiveDecryptedEnv).not.toHaveBeenCalled() + expect(mockGetEffectiveEnvironmentSnapshot).not.toHaveBeenCalled() expect(mockExecuteInternalToolOperation).not.toHaveBeenCalled() }) @@ -5258,7 +5320,7 @@ describe('Copilot Env Variable Reference Resolution', () => { expect(result.success).toBe(false) expect(result.error).toContain('only personal variables are available') - expect(mockGetEffectiveDecryptedEnv).toHaveBeenCalledWith('user-123', undefined) + expect(mockGetEffectiveEnvironmentSnapshot).toHaveBeenCalledWith('user-123', undefined) expect(mockExecuteInternalToolOperation).not.toHaveBeenCalled() }) @@ -5270,7 +5332,7 @@ describe('Copilot Env Variable Reference Resolution', () => { ) expect(result.success).toBe(true) - expect(mockGetEffectiveDecryptedEnv).not.toHaveBeenCalled() + expect(mockGetEffectiveEnvironmentSnapshot).not.toHaveBeenCalled() expect(sentOperationInput().apiKey).toBe('{{SENTRY_AUTH_TOKEN}}') }) diff --git a/apps/sim/tools/index.ts b/apps/sim/tools/index.ts index a3123e848be..748c4d8a7c7 100644 --- a/apps/sim/tools/index.ts +++ b/apps/sim/tools/index.ts @@ -471,8 +471,10 @@ async function resolveToolEnvReferences( const completePendingActivation = resolvedSecretTraceRegistry?.beginPendingActivation() try { - const { getEffectiveDecryptedEnv } = await import('@/lib/environment/utils') - const envVars = await getEffectiveDecryptedEnv(scope.userId, scope.workspaceId) + const environmentScope = { userId: scope.userId, workspaceId: scope.workspaceId } + const { getEffectiveEnvironmentSnapshot } = await import('@/lib/environment/utils') + const environment = await getEffectiveEnvironmentSnapshot(scope.userId, scope.workspaceId) + const envVars = { ...environment.personalDecrypted, ...environment.workspaceDecrypted } for (const { paramId, value, soft } of pending) { const missingKeys: string[] = [] @@ -480,9 +482,12 @@ async function resolveToolEnvReferences( allowEmbedded: false, missingKeys, onResolved: (name, resolvedValue) => { - resolvedSecretTraceRegistry?.recordResolvedAtInputPath(name, resolvedValue, [paramId], { - propagated: true, - }) + resolvedSecretTraceRegistry?.recordResolvedFromEnvironment( + name, + resolvedValue, + { ...environment, scope: environmentScope }, + { path: [paramId], propagated: true } + ) }, }) if (missingKeys.length > 0) {