diff --git a/apps/sim/lib/execution/durable-secret-provenance-telemetry.ts b/apps/sim/lib/execution/durable-secret-provenance-telemetry.ts index bf77d2be246..59d8ef72e4e 100644 --- a/apps/sim/lib/execution/durable-secret-provenance-telemetry.ts +++ b/apps/sim/lib/execution/durable-secret-provenance-telemetry.ts @@ -1,3 +1,10 @@ +import { + AuditAction, + type AuditLogParams, + AuditResourceType, + recordAudit, + recordAuditBatch, +} from '@sim/audit' import { createLogger } from '@sim/logger' const persistenceLogger = createLogger('DurableSecretProvenancePersistence') @@ -28,7 +35,6 @@ export type DurableSecretProvenanceRefusalCause = | 'workspace-file-provenance-unavailable' | 'workspace-file-opaque-secret-content' | 'workspace-file-registry-unavailable' - | 'workspace-file-unrecorded-enforced' export interface DurableSecretProvenanceRefusalReport { surface: DurableSecretProvenanceSurface @@ -65,3 +71,56 @@ export function reportDurableSecretProvenanceRefusal( ...(report.resourceId ? { resourceId: report.resourceId } : {}), }) } + +export interface DurableSecretProvenanceUnrecordedReport { + surface: DurableSecretProvenanceSurface + workspaceId?: string + organizationId?: string + resourceId?: string + recordCount?: number + actorUserId?: string +} + +function unrecordedProvenanceAuditEntry( + report: DurableSecretProvenanceUnrecordedReport +): AuditLogParams | undefined { + const metadata = { + surface: report.surface, + ...(report.organizationId ? { organizationId: report.organizationId } : {}), + ...(report.recordCount !== undefined ? { recordCount: report.recordCount } : {}), + } + persistenceLogger.warn('Using content without recorded secret provenance', { + ...metadata, + ...(report.workspaceId ? { workspaceId: report.workspaceId } : {}), + ...(report.resourceId ? { resourceId: report.resourceId } : {}), + }) + if (!report.workspaceId && !report.organizationId) return undefined + return { + workspaceId: report.organizationId ? null : report.workspaceId, + actorId: report.actorUserId ?? null, + action: AuditAction.SECRET_PROVENANCE_UNRECORDED, + resourceType: AuditResourceType.SECRET_PROVENANCE, + ...(report.resourceId ? { resourceId: report.resourceId } : {}), + description: 'Used content without recorded secret provenance', + metadata, + } +} + +/** Records accepted content whose producer did not supply provenance, without recording bytes. */ +export function reportDurableSecretProvenanceUnrecorded( + report: DurableSecretProvenanceUnrecordedReport +): void { + const entry = unrecordedProvenanceAuditEntry(report) + if (entry) recordAudit(entry) +} + +/** Batches tenant-scoped admission events without one pooled database query per source. */ +export function reportDurableSecretProvenanceUnrecordedBatch( + reports: readonly DurableSecretProvenanceUnrecordedReport[] +): void { + const entries = reports.flatMap((report) => { + const entry = unrecordedProvenanceAuditEntry(report) + return entry ? [entry] : [] + }) + if (entries.length > 0) recordAuditBatch(entries) +} diff --git a/apps/sim/lib/function-execution/execute-request.test.ts b/apps/sim/lib/function-execution/execute-request.test.ts index 548785d80ba..dd1e8e478b6 100644 --- a/apps/sim/lib/function-execution/execute-request.test.ts +++ b/apps/sim/lib/function-execution/execute-request.test.ts @@ -2264,11 +2264,13 @@ describe('Function execution request', () => { ) it.each([ - { reason: 'source-provenance-incomplete', status: 200 }, - { reason: 'entry-decrypt-failed', status: 400 }, + { reason: 'source-provenance-incomplete', status: 200, knownMount: false, text: false }, + { reason: 'source-provenance-incomplete', status: 200, knownMount: true, text: false }, + { reason: 'source-provenance-incomplete', status: 200, knownMount: true, text: true }, + { reason: 'entry-decrypt-failed', status: 400, knownMount: false, text: false }, ] as const)( - 'distinguishes historical absence from provenance faults: $reason', - async ({ reason, status }) => { + 'distinguishes historical absence from provenance faults: $reason knownMount=$knownMount text=$text', + async ({ reason, status, knownMount, text }) => { envFlagsMock.isRemoteSandboxEnabled = true const registry = new ResolvedSecretTraceRegistry([], { userId: 'user-123', @@ -2278,7 +2280,7 @@ describe('Function execution request', () => { const contentUpdatedAt = new Date('2026-01-01T00:00:00Z') mockMountContributors.mockReturnValue([ { - fileId: 'legacy-file', + fileId: knownMount ? 'known-file' : 'legacy-file', key: 'execution/workspace-1/workflow-1/execution-1/a/input.txt', context: 'execution', contentUpdatedAt, @@ -2287,25 +2289,38 @@ describe('Function execution request', () => { dbChainMockFns.limit.mockResolvedValue([ { fileContentUpdatedAt: contentUpdatedAt, - secretProvenanceVersion: null, - provenanceContentUpdatedAt: null, - status: null, - entries: null, + secretProvenanceVersion: knownMount ? 1 : null, + provenanceContentUpdatedAt: knownMount ? contentUpdatedAt : null, + status: knownMount ? 'exact' : null, + entries: knownMount + ? [ + { + name: 'API_KEY', + encryptedValue: 'encrypted:mounted-secret', + sourceUserId: 'user-123', + sourceWorkspaceId: 'workspace-1', + }, + ] + : null, }, ]) - const buffer = Buffer.from('ordinary file') - mockExecuteInSandbox.mockResolvedValueOnce({ + const buffer = Buffer.from(text ? 'Bearer mounted-secret' : 'ordinary file') + mockExecuteInSandbox.mockResolvedValue({ result: null, stdout: '', sandboxId: 'sbx', - collectedFiles: [ - { - relativePath: 'report.zip', - path: '/tmp/sim/outputs/report.zip', - contentBase64: buffer.toString('base64'), - byteLength: buffer.length, - }, - ], + ...(text + ? { exportedFiles: { '/home/user/report.txt': buffer.toString('utf8') } } + : { + collectedFiles: [ + { + relativePath: 'report.zip', + path: '/tmp/sim/outputs/report.zip', + contentBase64: buffer.toString('base64'), + byteLength: buffer.length, + }, + ], + }), }) const response = await POST( createMockRequest('POST', { @@ -2314,13 +2329,42 @@ describe('Function execution request', () => { workspaceId: 'workspace-1', workflowId: 'workflow-1', executionId: 'execution-1', + ...(text + ? { + outputs: { + files: [ + { + path: 'files/report.txt', + sandboxPath: '/home/user/report.txt', + mimeType: 'text/plain', + }, + ], + }, + } + : {}), }), registry ) expect(response.status).toBe(status) - if (status === 200) - expect(mockUploadExecutionFile.mock.calls[0][5]).toEqual({ status: 'unrecorded' }) - else expect(mockUploadExecutionFile).not.toHaveBeenCalled() + if (status === 200) { + if (text) { + expect(mockWriteWorkspaceFileByPath.mock.calls[0][0].secretProvenance).toEqual({ + status: 'exact', + entries: [ + { + name: 'API_KEY', + encryptedValue: 'encrypted:mounted-secret', + sourceUserId: 'user-123', + sourceWorkspaceId: 'workspace-1', + }, + ], + }) + } else { + expect(mockUploadExecutionFile.mock.calls[0][5]).toEqual({ + status: knownMount ? 'unknown' : 'unrecorded', + }) + } + } else expect(mockUploadExecutionFile).not.toHaveBeenCalled() } ) @@ -2567,7 +2611,7 @@ describe('Function execution request', () => { it.each([ ['a mount with no provenance source', true], ['no mounts', false], - ] as const)('withholds workbench certification for %s', async (_label, mounted) => { + ] as const)('keeps workbench use available for %s', async (_label, mounted) => { envFlagsMock.isMothershipSandboxEnabled = true mockUnprovenancedMountCount.mockReturnValue(mounted ? 1 : 0) hybridAuthMockFns.mockCheckInternalAuth.mockResolvedValue({ @@ -2588,7 +2632,7 @@ describe('Function execution request', () => { expect(response.status).toBe(200) const session = mockExecuteInSandbox.mock.calls.at(-1)?.[0].session expect(session.key).toBe('chat-session') - expect(session.unprovenancedInputs === true).toBe(mounted) + expect(session.unprovenancedInputs).not.toBe(true) }) it('gives overlapping calls in one persistent workbench distinct automatic export directories', async () => { diff --git a/apps/sim/lib/function-execution/execute-request.ts b/apps/sim/lib/function-execution/execute-request.ts index 143e14f642b..69a74ea8737 100644 --- a/apps/sim/lib/function-execution/execute-request.ts +++ b/apps/sim/lib/function-execution/execute-request.ts @@ -1089,6 +1089,10 @@ async function importRuntimeInputProvenance( }) if (decision.safe && decision.provenance.status === 'unrecorded') { context.runtimeInputProvenanceUnrecorded = true + context.runtimeFileSecretTraceRegistry = new ResolvedSecretTraceRegistry([], { + userId: context.attributedUserId, + workspaceId: context.workspaceId, + }) return } } @@ -2718,6 +2722,9 @@ export async function executeFunctionRequest( logger, }, }) + if (resolvedMounts.unprovenancedMountCount > 0) { + routeContext.runtimeInputProvenanceUnrecorded = true + } await importRuntimeFileContributors( routeContext, resolvedMounts.contributingFiles, @@ -2747,7 +2754,6 @@ export async function executeFunctionRequest( const mothershipSession = admittedSession ? { ...admittedSession, - unprovenancedInputs: resolvedMounts.unprovenancedMountCount > 0, inputProvenance: () => { const runtime = activeRouteContext.runtimeFileSecretTraceRegistry?.exportProvenance() return mergeDurableSecretProvenance( diff --git a/apps/sim/lib/function-execution/sandbox-mounts.ts b/apps/sim/lib/function-execution/sandbox-mounts.ts index 6c9faf28557..9fb3f3c836d 100644 --- a/apps/sim/lib/function-execution/sandbox-mounts.ts +++ b/apps/sim/lib/function-execution/sandbox-mounts.ts @@ -1,4 +1,5 @@ import { createLogger } from '@sim/logger' +import { reportDurableSecretProvenanceUnrecorded } from '@/lib/execution/durable-secret-provenance-telemetry' import { resolveStoredFileProvenanceSource } from '@/lib/execution/payloads/file-secret-provenance' import { assertUserFileContentAccess, @@ -277,15 +278,7 @@ export async function resolveUserFileMounts(args: { manifest: SandboxMountManifestEntry[] contributingFiles?: readonly WorkspaceFileSecretProvenanceIdentity[] renderedContributingFiles?: readonly WorkspaceFileSecretProvenanceIdentity[] - /** - * Mounts whose own bytes have no provenance source (no principal to bind one, or a key with no - * canonical metadata record). Workflow runs keep their legacy absence policy; a persistent - * workbench must not certify a machine that received one. - * - * Storage contexts other than workspace and execution (chat, copilot, knowledge-base, logs, and - * the other public contexts) never have a source, so they always count here and taint a - * workbench. That is conservative by design. - */ + /** Mounts without producer provenance remain usable and report the recording gap. */ unprovenancedMountCount: number }> { let unprovenancedMountCount = 0 @@ -365,6 +358,14 @@ export async function resolveUserFileMounts(args: { }) } + if (unprovenancedMountCount > 0) { + reportDurableSecretProvenanceUnrecorded({ + surface: 'workspace-file', + workspaceId: args.context.workspaceId, + actorUserId: args.context.userId, + recordCount: unprovenancedMountCount, + }) + } logger.info('Resolved sandbox file mounts', { mountCount: sandboxFiles.length, bufferedBytes: budget.buffered, diff --git a/apps/sim/lib/internal/file/operations.provenance.test.ts b/apps/sim/lib/internal/file/operations.provenance.test.ts index 19a7f2a0dc4..fda49a84a38 100644 --- a/apps/sim/lib/internal/file/operations.provenance.test.ts +++ b/apps/sim/lib/internal/file/operations.provenance.test.ts @@ -89,6 +89,8 @@ vi.mock('@/lib/workspaces/permissions/utils', () => permissionsMock) vi.mock('@/app/api/files/authorization', () => filesAuthorizationMock) vi.mock('@/lib/execution/durable-secret-provenance-telemetry', () => ({ + reportDurableSecretProvenanceUnrecorded: vi.fn(), + reportDurableSecretProvenanceUnrecordedBatch: vi.fn(), reportDurableSecretProvenanceWrite: vi.fn(), reportDurableSecretProvenanceRefusal: vi.fn(), })) @@ -237,7 +239,7 @@ describe('appended file provenance', () => { }) it.each([ - { predecessor: 'unrecorded', secret: true, expectedStatus: 'unknown' }, + { predecessor: 'unrecorded', secret: true, expectedStatus: 'exact' }, { predecessor: 'exact', secret: true, expectedStatus: 'exact' }, { predecessor: 'legacy', secret: true, expectedStatus: 'exact' }, { predecessor: 'unrecorded', secret: false, expectedStatus: 'unrecorded' }, @@ -315,7 +317,7 @@ describe('appended file provenance', () => { view: 'complete', value: `before:${content}`, }) - expect(permitted).toBe(expectedStatus === 'exact') + expect(permitted).toBe(expectedStatus !== 'unknown') if (permitted) { expect(projectResolvedSecretModelContent(`before:${content}`, registry)).toEqual({ safe: true, @@ -326,7 +328,7 @@ describe('appended file provenance', () => { joinedRow(persisted.status, persisted.entries, persisted.contentUpdatedAt), ]) expect(await isOpaqueWorkspaceFileEgressSafe('workspace-1', IDENTITY)).toBe( - expectedStatus === 'exact' && !secret + expectedStatus !== 'unknown' && !secret ) } ) @@ -352,7 +354,7 @@ describe('execution-file content provenance', () => { it.each([ { status: 'exact', version: 1, stale: false, complete: true }, - { status: 'unrecorded', version: 1, stale: false, complete: false }, + { status: 'unrecorded', version: 1, stale: false, complete: true }, { status: 'unknown', version: 1, stale: false, complete: false }, { status: 'unknown', version: null, stale: false, complete: true }, { status: 'exact', version: 1, stale: true, complete: false }, diff --git a/apps/sim/lib/internal/file/operations.ts b/apps/sim/lib/internal/file/operations.ts index 1ac6587c30f..eebc28d1c07 100644 --- a/apps/sim/lib/internal/file/operations.ts +++ b/apps/sim/lib/internal/file/operations.ts @@ -14,6 +14,7 @@ import { isPayloadSizeLimitError } from '@/lib/core/utils/stream-limits' import { ensureAbsoluteUrl } from '@/lib/core/utils/urls' import { isUserFile } from '@/lib/core/utils/user-file' import { durableSecretProvenanceFromPrivateBundle } from '@/lib/execution/durable-secret-provenance' +import { reportDurableSecretProvenanceUnrecorded } from '@/lib/execution/durable-secret-provenance-telemetry' import { inspectPrivateSecretProvenanceRequest, isPrivateSecretProvenanceBundleV1, @@ -596,6 +597,7 @@ export async function getFileContentProvenance( ? { userId: ownerUserId, workspaceId } : undefined const accumulator = new ResolvedSecretTraceProvenanceAccumulator(scope) + let unrecorded = 0 for (const source of sources) { signal?.throwIfAborted() @@ -605,7 +607,11 @@ export async function getFileContentProvenance( } const provenance = await readFileSourceSecretProvenance(principal, workspaceId, source.identity) signal?.throwIfAborted() - if (provenance.status !== 'exact' || (source.opaque && provenance.entries.length > 0)) { + if (provenance.status === 'unrecorded') { + unrecorded += 1 + continue + } + if (provenance.status === 'unknown' || (source.opaque && provenance.entries.length > 0)) { accumulator.markIncomplete('workspace-file-provenance-unknown') continue } @@ -625,6 +631,13 @@ export async function getFileContentProvenance( }) } + if (unrecorded > 0) { + reportDurableSecretProvenanceUnrecorded({ + surface: 'workspace-file', + workspaceId, + recordCount: unrecorded, + }) + } return accumulator.exportProvenance() } diff --git a/apps/sim/lib/knowledge/__integration__/execution-archive-provenance.integration.ts b/apps/sim/lib/knowledge/__integration__/execution-archive-provenance.integration.ts index f2e85e1983d..3a3332a2569 100644 --- a/apps/sim/lib/knowledge/__integration__/execution-archive-provenance.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/execution-archive-provenance.integration.ts @@ -1,4 +1,4 @@ -/** Real execution-file storage, ZIP extraction, durable provenance, table import, and KB indexing. */ +/** Real execution-file storage, ZIP extraction, durable provenance, and KB indexing. */ import { mkdtempSync } from 'node:fs' import { rm } from 'node:fs/promises' import { tmpdir } from 'node:os' @@ -182,70 +182,74 @@ afterAll(async () => { }) describe('execution archive durable provenance', () => { - it('carries exact-empty lineage through extraction and delayed KB indexing/search', async () => { - const ids = await seed() - const archive = await uploadArchive(ids, { status: 'exact', entries: [] }) - const [storedArchive] = await db - .select({ - secretProvenanceVersion: workspaceFiles.secretProvenanceVersion, - context: workspaceFiles.context, - }) - .from(workspaceFiles) - .where(eq(workspaceFiles.key, archive.key)) - expect(storedArchive.secretProvenanceVersion).toBe(1) - expect(storedArchive.context).toBe('execution') - const source = await extract(ids, archive) - expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, source.identity)).toEqual({ - status: 'exact', - entries: [], - }) - expect(await isOpaqueWorkspaceFileEgressSafe(ids.workspaceId, source.identity)).toBe(true) - expect((await downloadFile({ key: source.child.key, context: 'workspace' })).toString()).toBe( - REPORT_CSV - ) + it.each(['exact', 'unrecorded'] as const)( + 'carries %s lineage through extraction and delayed KB indexing/search', + async (status) => { + const ids = await seed() + const provenance: WorkspaceFileSecretProvenance = + status === 'exact' ? { status, entries: [] } : { status } + const archive = await uploadArchive(ids, provenance) + const [storedArchive] = await db + .select({ + secretProvenanceVersion: workspaceFiles.secretProvenanceVersion, + context: workspaceFiles.context, + }) + .from(workspaceFiles) + .where(eq(workspaceFiles.key, archive.key)) + expect(storedArchive.secretProvenanceVersion).toBe(1) + expect(storedArchive.context).toBe('execution') + const source = await extract(ids, archive) + expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, source.identity)).toEqual( + provenance + ) + expect(await isOpaqueWorkspaceFileEgressSafe(ids.workspaceId, source.identity)).toBe(true) + expect((await downloadFile({ key: source.child.key, context: 'workspace' })).toString()).toBe( + REPORT_CSV + ) - const imported = await addWorkspaceFilesToKnowledgeBase.execute({ - principal: sessionPrincipal(ids), - input: { knowledgeBaseId: ids.knowledgeBaseId, fileReferences: [source.child.id] }, - }) - expect(imported.failed).toEqual([]) - expect(imported.added).toHaveLength(1) - const documentId = imported.added[0].documentId - const [admitted] = await db.select().from(document).where(eq(document.id, documentId)) - expect(admitted.secretProvenanceVersion).toBe(1) - expect(admitted.storageKey).toMatch(/^kb\//) - const events = await db - .select() - .from(outboxEvent) - .where(sql`${outboxEvent.payload}::jsonb ->> 'documentId' = ${documentId}`) - trackedEventIds.push(...events.map((event) => event.id)) - const dispatch = events.find( - (event) => event.eventType === KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT - ) - if (!dispatch) throw new Error('Knowledge import did not atomically admit processing') - await deleteWorkspaceFile(ids.workspaceId, source.child.id) - await deleteFile({ key: source.child.key, context: 'workspace' }) - await deleteFile({ key: archive.key, context: 'execution' }) - await processOutboxEventById(dispatch.id, knowledgeDocumentProcessingOutboxHandlers) - const [indexed] = await db.select().from(document).where(eq(document.id, documentId)) - expect(indexed.processingStatus, indexed.processingError ?? undefined).toBe('completed') - const chunks = await listKnowledgeChunks.execute({ - principal: sessionPrincipal(ids), - input: { knowledgeBaseId: ids.knowledgeBaseId, documentId }, - }) - expect(chunks.chunks.map((chunk) => chunk.content).join('\n')).toContain(REPORT_TEXT) - const search = await searchKnowledge.execute({ - principal: sessionPrincipal(ids), - input: { - workspaceId: ids.workspaceId, - knowledgeBaseIds: [ids.knowledgeBaseId], - query: 'Orion', - searchMode: 'hybrid', - topK: 10, - }, - }) - expect(search.results.map((entry) => entry.documentId)).toContain(documentId) - }) + const imported = await addWorkspaceFilesToKnowledgeBase.execute({ + principal: sessionPrincipal(ids), + input: { knowledgeBaseId: ids.knowledgeBaseId, fileReferences: [source.child.id] }, + }) + expect(imported.failed).toEqual([]) + expect(imported.added).toHaveLength(1) + const documentId = imported.added[0].documentId + const [admitted] = await db.select().from(document).where(eq(document.id, documentId)) + expect(admitted.secretProvenanceVersion).toBe(1) + expect(admitted.storageKey).toMatch(/^kb\//) + const events = await db + .select() + .from(outboxEvent) + .where(sql`${outboxEvent.payload}::jsonb ->> 'documentId' = ${documentId}`) + trackedEventIds.push(...events.map((event) => event.id)) + const dispatch = events.find( + (event) => event.eventType === KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT + ) + if (!dispatch) throw new Error('Knowledge import did not atomically admit processing') + await deleteWorkspaceFile(ids.workspaceId, source.child.id) + await deleteFile({ key: source.child.key, context: 'workspace' }) + await deleteFile({ key: archive.key, context: 'execution' }) + await processOutboxEventById(dispatch.id, knowledgeDocumentProcessingOutboxHandlers) + const [indexed] = await db.select().from(document).where(eq(document.id, documentId)) + expect(indexed.processingStatus, indexed.processingError ?? undefined).toBe('completed') + const chunks = await listKnowledgeChunks.execute({ + principal: sessionPrincipal(ids), + input: { knowledgeBaseId: ids.knowledgeBaseId, documentId }, + }) + expect(chunks.chunks.map((chunk) => chunk.content).join('\n')).toContain(REPORT_TEXT) + const search = await searchKnowledge.execute({ + principal: sessionPrincipal(ids), + input: { + workspaceId: ids.workspaceId, + knowledgeBaseIds: [ids.knowledgeBaseId], + query: 'Orion', + searchMode: 'hybrid', + topK: 10, + }, + }) + expect(search.results.map((entry) => entry.documentId)).toContain(documentId) + } + ) it('keeps an explicitly unknown execution source unavailable to model and KB consumers', async () => { const ids = await seed() @@ -345,8 +349,10 @@ describe('execution archive durable provenance', () => { it.each([ { status: 'exact', deleted: false }, { status: 'unknown', deleted: false }, + { status: 'unrecorded', deleted: false }, { status: 'exact', deleted: true }, { status: 'unknown', deleted: true }, + { status: 'unrecorded', deleted: true }, ] as const)( 'binds $status execution bytes into KB admission despite URL-only classification (deleted=$deleted)', async ({ status, deleted }) => { @@ -390,12 +396,12 @@ describe('execution archive durable provenance', () => { .from(document) .leftJoin(documentSecretProvenance, eq(documentSecretProvenance.documentId, document.id)) .where(eq(document.id, admitted.id)) - expect(stored).toEqual({ version: 1, status }) + expect(stored).toEqual({ version: 1, status: status === 'unrecorded' ? 'exact' : status }) const registry = loadKnowledgeDocumentSecretRegistry(admitted.id, { userId: ids.aliceId, workspaceId: ids.workspaceId, }) - if (status === 'exact') { + if (status !== 'unknown') { await expect(registry).resolves.toMatchObject({ tracked: true, provenance: { status: 'exact', entries: [] }, diff --git a/apps/sim/lib/knowledge/__integration__/external-file-provenance.integration.ts b/apps/sim/lib/knowledge/__integration__/external-file-provenance.integration.ts index e4bd466e9e9..39c7192f9d2 100644 --- a/apps/sim/lib/knowledge/__integration__/external-file-provenance.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/external-file-provenance.integration.ts @@ -167,19 +167,19 @@ describe('external file ingress under durable enforcement', () => { }) ).toBe(true) - /** A real unrecorded control proves this suite has not silently disabled enforcement. */ - const unrecorded = await uploadExecutionFile( + /** A real unknown control proves missing evidence has not disabled protection failures. */ + const unknown = await uploadExecutionFile( ids, IMAGE_BYTES, - 'unrecorded.png', + 'unknown.png', 'image/png', ids.aliceId, - { status: 'unrecorded' } + { status: 'unknown' } ) expect( await importWorkspaceFileSecretProvenanceForModelView({ workspaceId: ids.workspaceId, - identity: await identityFor(unrecorded), + identity: await identityFor(unknown), 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 index 1b7f4252fa1..16733d0e868 100644 --- a/apps/sim/lib/knowledge/__integration__/upload-read-provenance.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/upload-read-provenance.integration.ts @@ -3,17 +3,20 @@ import { mkdtempSync } from 'node:fs' import { rm } from 'node:fs/promises' import { tmpdir } from 'node:os' import path from 'node:path' +import { AuditAction } from '@sim/audit' import { db } from '@sim/db' import { + auditLog, copilotChats, knowledgeBase, organization, user, workspace, + workspaceFileSecretProvenance, workspaceFiles, } from '@sim/db/schema' import { generateId } from '@sim/utils/id' -import { eq, inArray } from 'drizzle-orm' +import { and, eq, inArray } from 'drizzle-orm' import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest' const fixtureStorage = vi.hoisted(() => ({ root: '' })) @@ -32,9 +35,16 @@ import { trackChatUpload, } from '@/lib/uploads/contexts/workspace/workspace-file-manager' import { + areModelSafeWorkspaceFileKeys, + filterModelSafeWorkspaceFileAttachments, getBoundWorkspaceFileSecretProvenance, + getBoundWorkspaceFileSecretProvenanceByMetadata, importWorkspaceFileSecretProvenanceForModelView, + importWorkspaceFileSecretProvenanceForRuntime, + isModelSafeWorkspaceFileKey, + isOpaqueWorkspaceFileEgressSafe, replaceWorkspaceFileSecretProvenanceInTx, + snapshotWorkspaceFileSecretProvenanceInTx, type WorkspaceFileSecretProvenance, } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' import { uploadFile } from '@/lib/uploads/core/storage-service' @@ -72,7 +82,7 @@ async function seedUpload(provenance?: WorkspaceFileSecretProvenance) { key, name, 'text/plain', - CONTENT.length + Buffer.byteLength(CONTENT) ) const [file] = await db.select().from(workspaceFiles).where(eq(workspaceFiles.key, key)) if (provenance) { @@ -118,6 +128,7 @@ beforeAll(() => { afterAll(async () => { for (const ids of fixtures) { + await db.delete(auditLog).where(eq(auditLog.workspaceId, ids.workspaceId)) 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)) @@ -128,10 +139,16 @@ afterAll(async () => { }) describe('chat upload reads racing with save_upload', () => { - it.each(['legacy', 'exact'] as const)( + it.each(['legacy', 'exact', 'unrecorded'] as const)( 'keeps a %s authorized upload read valid across promotion', async (kind) => { - const ids = await seedUpload(kind === 'exact' ? { status: 'exact', entries: [] } : undefined) + const ids = await seedUpload( + kind === 'legacy' + ? undefined + : kind === 'exact' + ? { status: 'exact', entries: [] } + : { status: 'unrecorded' } + ) const read = await readUpload(ids) expect(read?.value.content).toBe(CONTENT) expect(read?.file).toBeDefined() @@ -217,3 +234,118 @@ describe('chat upload reads racing with save_upload', () => { } ) }) + +describe('recorded file absence policy', () => { + it.each(['keys', 'attachments'] as const)( + 'records unscoped admitted %s once per canonical workspace', + async (boundary) => { + const first = await seedUpload({ status: 'unrecorded' }) + const second = await seedUpload({ status: 'unrecorded' }) + const sibling = { + ...first.file, + id: generateId(), + key: `${first.file.key}-sibling`, + originalName: 'sibling.txt', + displayName: 'sibling.txt', + } + await db.insert(workspaceFiles).values(sibling) + await db.transaction((tx) => + replaceWorkspaceFileSecretProvenanceInTx(tx, sibling.id, sibling.contentUpdatedAt, { + status: 'unrecorded', + }) + ) + const files = [first.file, sibling, second.file] + if (boundary === 'keys') { + expect(await areModelSafeWorkspaceFileKeys(files.map((file) => file.key))).toBe(true) + } else { + const attachments = files.map((file) => ({ id: file.id, key: file.key })) + expect(await filterModelSafeWorkspaceFileAttachments(attachments)).toEqual(attachments) + } + const query = () => + db + .select({ workspaceId: auditLog.workspaceId, metadata: auditLog.metadata }) + .from(auditLog) + .where( + and( + inArray(auditLog.workspaceId, [first.workspaceId, second.workspaceId]), + eq(auditLog.action, AuditAction.SECRET_PROVENANCE_UNRECORDED) + ) + ) + await expect.poll(query).toHaveLength(2) + const rows = await query() + expect(rows.find((row) => row.workspaceId === first.workspaceId)?.metadata).toMatchObject({ + recordCount: 2, + }) + expect(rows.find((row) => row.workspaceId === second.workspaceId)?.metadata).toMatchObject({ + recordCount: 1, + }) + if (boundary === 'keys') { + await db + .delete(auditLog) + .where(inArray(auditLog.workspaceId, [first.workspaceId, second.workspaceId])) + await db.transaction((tx) => + replaceWorkspaceFileSecretProvenanceInTx(tx, sibling.id, sibling.contentUpdatedAt, { + status: 'unknown', + }) + ) + expect(await areModelSafeWorkspaceFileKeys(files.map((file) => file.key))).toBe(false) + expect(await query()).toEqual([]) + } + } + ) + + it('keeps valid unrecorded bytes usable across runtime, opaque and attachment boundaries', async () => { + const ids = await seedUpload({ status: 'unrecorded' }) + const read = await readUpload(ids) + expect( + await importWorkspaceFileSecretProvenanceForRuntime({ + workspaceId: ids.workspaceId, + identity: read.file, + }) + ).toBe(true) + expect(await isOpaqueWorkspaceFileEgressSafe(ids.workspaceId, read.file)).toBe(true) + expect(await isModelSafeWorkspaceFileKey(ids.file.key, { workspaceId: ids.workspaceId })).toBe( + true + ) + const attachment = { id: ids.file.id, key: ids.file.key } + expect( + await filterModelSafeWorkspaceFileAttachments([attachment], { workspaceId: ids.workspaceId }) + ).toEqual([attachment]) + expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, read.file)).toEqual({ + status: 'unrecorded', + }) + }) + + it('does not discard recorded entries from an inconsistent unrecorded sidecar', async () => { + const ids = await seedUpload({ status: 'unrecorded' }) + await db + .update(workspaceFileSecretProvenance) + .set({ + entries: [ + { name: 'TOKEN', encryptedValue: 'fixture-ciphertext', sourceUserId: ids.aliceId }, + ], + }) + .where(eq(workspaceFileSecretProvenance.fileId, ids.file.id)) + const read = await readUpload(ids) + expect(await getBoundWorkspaceFileSecretProvenance(ids.workspaceId, read.file)).toEqual({ + status: 'unknown', + }) + const batch = await getBoundWorkspaceFileSecretProvenanceByMetadata(db, [ + { ...ids.file, secretProvenanceVersion: 1 }, + ]) + expect(batch.get(ids.file.id)).toEqual({ status: 'unknown' }) + expect( + await db.transaction((tx) => + snapshotWorkspaceFileSecretProvenanceInTx(tx, ids.file.id, ids.file.contentUpdatedAt, 1) + ) + ).toEqual({ status: 'unknown', entries: [] }) + expect(await isModelSafeWorkspaceFileKey(ids.file.key, { workspaceId: ids.workspaceId })).toBe( + false + ) + expect( + await filterModelSafeWorkspaceFileAttachments([{ id: ids.file.id, key: ids.file.key }], { + workspaceId: ids.workspaceId, + }) + ).toEqual([]) + }) +}) diff --git a/apps/sim/lib/knowledge/__integration__/workspace-import.integration.ts b/apps/sim/lib/knowledge/__integration__/workspace-import.integration.ts index 55014eb1470..113eaa599e0 100644 --- a/apps/sim/lib/knowledge/__integration__/workspace-import.integration.ts +++ b/apps/sim/lib/knowledge/__integration__/workspace-import.integration.ts @@ -3,12 +3,16 @@ import { mkdtempSync } from 'node:fs' import { rm } from 'node:fs/promises' import { tmpdir } from 'node:os' import path from 'node:path' +import * as audit from '@sim/audit' +import { AuditAction } from '@sim/audit' import { db } from '@sim/db' import { + auditLog, document, knowledgeBase, organization, outboxEvent, + permissions, user, workspace, workspaceFiles, @@ -46,13 +50,14 @@ import { listKnowledgeChunks } from '@/lib/knowledge/application/chunks' import { searchKnowledge } from '@/lib/knowledge/application/search' import { KNOWLEDGE_DOCUMENT_PROCESSING_OUTBOX_EVENT } from '@/lib/knowledge/documents/processing-outbox-event' import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler' -import { createSingleDocument } from '@/lib/knowledge/documents/service' +import { createDocumentRecords, createSingleDocument } from '@/lib/knowledge/documents/service' import { KNOWLEDGE_STORAGE_CLEANUP_EVENT } from '@/lib/knowledge/documents/storage-cleanup' import { uploadKnowledgeArtifact } from '@/lib/knowledge/documents/storage-upload' import { deleteWorkspaceFile, uploadWorkspaceFile, } from '@/lib/uploads/contexts/workspace/workspace-file-manager' +import type { WorkspaceFileSecretProvenance } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' import { deleteFile, downloadFile } from '@/lib/uploads/core/storage-service' const fixtures: ReturnType[] = [] @@ -65,7 +70,11 @@ async function seed() { return ids } -async function sourceFile(ids: ReturnType, content: string) { +async function sourceFile( + ids: ReturnType, + content: string, + secretProvenance: WorkspaceFileSecretProvenance = { status: 'exact', entries: [] } +) { return uploadWorkspaceFile( ids.workspaceId, ids.aliceId, @@ -73,12 +82,44 @@ async function sourceFile(ids: ReturnType, 'import.txt', 'text/plain', { - secretProvenance: { status: 'exact', entries: [] }, + secretProvenance, notifyWorkspaceChange: false, } ) } +/** Drains real audit writes before asserting persisted rows, including duplicates from one operation. */ +function observeAbsenceAudits(workspaceId: string) { + let submitted = 0 + const recordAudit = audit.recordAudit + const observation = vi.spyOn(audit, 'recordAudit').mockImplementation((entry) => { + if ( + entry.workspaceId === workspaceId && + entry.action === AuditAction.SECRET_PROVENANCE_UNRECORDED + ) { + submitted += 1 + } + recordAudit(entry) + }) + return { + restore: () => observation.mockRestore(), + async persisted() { + const read = () => + db + .select() + .from(auditLog) + .where( + and( + eq(auditLog.workspaceId, workspaceId), + eq(auditLog.action, AuditAction.SECRET_PROVENANCE_UNRECORDED) + ) + ) + await vi.waitFor(async () => expect(await read()).toHaveLength(submitted)) + return read() + }, + } +} + beforeAll(() => { fixtureStorage.root = mkdtempSync(path.join(tmpdir(), 'sim-workspace-import-')) }) @@ -87,6 +128,7 @@ afterAll(async () => { if (trackedEventIds.length) await db.delete(outboxEvent).where(inArray(outboxEvent.id, trackedEventIds)) for (const ids of fixtures) { + await db.delete(auditLog).where(eq(auditLog.workspaceId, ids.workspaceId)) 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)) @@ -97,6 +139,96 @@ afterAll(async () => { }) describe('durable workspace file import', () => { + it('audits one accepted unrecorded source under the importer instead of its uploader', async () => { + const ids = await seed() + const source = await sourceFile(ids, 'ordinary source bytes', { status: 'unrecorded' }) + await db + .update(permissions) + .set({ permissionType: 'write' }) + .where(and(eq(permissions.userId, ids.bobId), eq(permissions.entityId, ids.workspaceId))) + const observed = observeAbsenceAudits(ids.workspaceId) + try { + const result = await addWorkspaceFilesToKnowledgeBase.execute({ + principal: { kind: 'session', userId: ids.bobId, sessionId: 'fixture-session' }, + input: { knowledgeBaseId: ids.knowledgeBaseId, fileReferences: [source.id, source.id] }, + }) + expect(result.failed).toEqual([]) + expect(result.added).toHaveLength(1) + const documentId = result.added[0].documentId + const events = await db + .select({ id: outboxEvent.id }) + .from(outboxEvent) + .where(sql`${outboxEvent.payload}::jsonb ->> 'documentId' = ${documentId}`) + trackedEventIds.push(...events.map((event) => event.id)) + const entries = await observed.persisted() + expect(entries).toHaveLength(1) + expect(entries[0]).toMatchObject({ + actorId: ids.bobId, + metadata: { surface: 'workspace-file', recordCount: 1 }, + }) + expect(await db.select().from(document).where(eq(document.id, documentId))).toHaveLength(1) + } finally { + observed.restore() + } + }) + + it.each([ + ['single', 'commit'], + ['single', 'rollback'], + ['bulk', 'commit'], + ['bulk', 'rollback'], + ] as const)('audits unrecorded %s document admission only after %s', async (mode, outcome) => { + const ids = await seed() + const source = await sourceFile(ids, 'ordinary source bytes', { status: 'unrecorded' }) + const observed = observeAbsenceAudits(ids.workspaceId) + const tripwire = process.env.DB_TX_TRIPWIRE + /** Match production's warning mode so a premature audit remains visible after rollback. */ + vi.stubEnv('DB_TX_TRIPWIRE', 'warn') + const transaction = db.transaction.bind(db) + const failure = + outcome === 'rollback' + ? vi.spyOn(db, 'transaction').mockImplementationOnce((callback, config) => + transaction(async (tx) => { + await callback(tx) + throw new Error('Document admission commit failed') + }, config) + ) + : undefined + try { + const input = { + filename: source.name, + fileUrl: `/api/files/serve/${encodeURIComponent(source.key)}?context=workspace`, + fileSize: source.size, + mimeType: source.type, + } + const created = + mode === 'single' + ? createSingleDocument(input, ids.knowledgeBaseId, generateId(), ids.aliceId) + : createDocumentRecords([input], ids.knowledgeBaseId, generateId(), ids.aliceId) + if (outcome === 'rollback') { + await expect(created).rejects.toThrow('Document admission commit failed') + } else { + await created + } + const expectedCount = outcome === 'commit' ? 1 : 0 + expect( + await db.select().from(document).where(eq(document.knowledgeBaseId, ids.knowledgeBaseId)) + ).toHaveLength(expectedCount) + const entries = await observed.persisted() + expect(entries).toHaveLength(expectedCount) + if (outcome === 'commit') { + expect(entries[0]).toMatchObject({ + actorId: null, + metadata: { surface: 'workspace-file', recordCount: 1 }, + }) + } + } finally { + failure?.mockRestore() + observed.restore() + vi.stubEnv('DB_TX_TRIPWIRE', tripwire) + } + }) + it('indexes the admitted snapshot after the source was deleted and temporary access would have expired', async () => { const ids = await seed() const content = diff --git a/apps/sim/lib/knowledge/application/add-workspace-files.test.ts b/apps/sim/lib/knowledge/application/add-workspace-files.test.ts index 97782b108c2..d149089ffa2 100644 --- a/apps/sim/lib/knowledge/application/add-workspace-files.test.ts +++ b/apps/sim/lib/knowledge/application/add-workspace-files.test.ts @@ -125,7 +125,10 @@ describe('add workspace files to knowledge base application command', () => { billedAccountUserId: 'billing-owner-1', }) workspaceFileSecretProvenanceMockFns.mockGetBoundWorkspaceFileSecretProvenance.mockResolvedValue( - { status: 'exact', entries: [] } + { + status: 'exact', + entries: [], + } ) workspaceFileManagerMockFns.mockFetchServableWorkspaceFileBuffer.mockResolvedValue({ buffer: Buffer.alloc(100), diff --git a/apps/sim/lib/knowledge/application/add-workspace-files.ts b/apps/sim/lib/knowledge/application/add-workspace-files.ts index b4822bf8a2f..241d463f469 100644 --- a/apps/sim/lib/knowledge/application/add-workspace-files.ts +++ b/apps/sim/lib/knowledge/application/add-workspace-files.ts @@ -1,11 +1,11 @@ import { AuditAction, AuditResourceType } from '@sim/audit' -import type { Principal } from '@sim/auth/principal' +import { type Principal, resolvePrincipalAuditAttribution } from '@sim/auth/principal' import { generateId } from '@sim/utils/id' import { checkAttributedUsageLimits } from '@/lib/billing/core/billing-attribution' import { authorizeWorkspaceOperation } from '@/lib/core/application' import { asOrchestrationError, OrchestrationError } from '@/lib/core/orchestration/types' import { generateRequestId } from '@/lib/core/utils/request' -import { reportDurableSecretProvenanceRefusal } from '@/lib/execution/durable-secret-provenance-telemetry' +import { reportDurableSecretProvenanceUnrecorded } from '@/lib/execution/durable-secret-provenance-telemetry' import { PROVENANCE_MAX_ENTRIES } from '@/lib/execution/provenance-limits' import { knowledgeDelegationPolicy } from '@/lib/knowledge/application/authorization' import { defineAuthorizedKnowledgeUseCase } from '@/lib/knowledge/application/authorized-knowledge-use-case' @@ -81,6 +81,7 @@ interface AddWorkspaceFilesContext extends ActiveKnowledgeBaseContext { interface PreparedWorkspaceFile { reference: string file: WorkspaceFileRecord + unrecordedProvenance: boolean } async function prepareWorkspaceFile( @@ -110,7 +111,7 @@ async function prepareWorkspaceFile( const fileTypeError = validateFileType(file.name, file.type) if (fileTypeError) throw new OrchestrationError('validation', fileTypeError.message) - await assertImportProvenance(context.workspaceId, { + const unrecordedProvenance = await assertImportProvenance(context.workspaceId, { fileId: file.id, key: file.key, context: file.storageContext ?? 'workspace', @@ -120,26 +121,25 @@ async function prepareWorkspaceFile( return { reference, file, + unrecordedProvenance, } } async function assertImportProvenance( workspaceId: string, identity: WorkspaceFileSecretProvenanceIdentity -): Promise { +): Promise { const provenance = await getBoundWorkspaceFileSecretProvenance(workspaceId, identity) - if (provenance.status !== 'exact' || provenance.entries.length > 0) { - reportDurableSecretProvenanceRefusal({ - surface: 'knowledge', - cause: 'knowledge-workspace-file-source-unavailable', - workspaceId, - resourceId: identity.fileId, - }) + if ( + provenance.status === 'unknown' || + (provenance.status === 'exact' && provenance.entries.length > 0) + ) { throw new OrchestrationError( 'validation', 'Workspace file secret provenance prevents knowledge ingestion' ) } + return provenance.status === 'unrecorded' } export const addWorkspaceFilesToKnowledgeBase = defineAuthorizedKnowledgeUseCase({ @@ -220,6 +220,8 @@ export const addWorkspaceFilesToKnowledgeBase = defineAuthorizedKnowledgeUseCase try { const requestId = generateRequestId() const current = await prepareWorkspaceFile(principal, context, candidate.file.id, 'id') + const unrecordedSources = new Set() + if (current.unrecordedProvenance) unrecordedSources.add(current.file.id) const readSignal = AbortSignal.timeout(SOURCE_READ_TIMEOUT_MS) const signal = input.cancellationSignal ? AbortSignal.any([readSignal, input.cancellationSignal]) @@ -242,7 +244,9 @@ export const addWorkspaceFilesToKnowledgeBase = defineAuthorizedKnowledgeUseCase } for (const identity of contributors) { signal.throwIfAborted() - await assertImportProvenance(context.workspaceId, identity) + if (await assertImportProvenance(context.workspaceId, identity)) { + unrecordedSources.add(identity.fileId) + } } const documentId = generateId() const storedFile = await uploadKnowledgeArtifact({ @@ -270,6 +274,7 @@ export const addWorkspaceFilesToKnowledgeBase = defineAuthorizedKnowledgeUseCase ) { throw new OrchestrationError('conflict', 'Workspace file changed during knowledge import') } + if (finalized.unrecordedProvenance) unrecordedSources.add(finalized.file.id) input.cancellationSignal?.throwIfAborted() const document = await createSingleDocument( { @@ -293,6 +298,14 @@ export const addWorkspaceFilesToKnowledgeBase = defineAuthorizedKnowledgeUseCase processing: { processingOptions: {}, billingAttribution }, } ) + if (unrecordedSources.size > 0) { + reportDurableSecretProvenanceUnrecorded({ + surface: 'workspace-file', + workspaceId: context.workspaceId, + actorUserId: resolvePrincipalAuditAttribution(principal).actorId ?? undefined, + recordCount: unrecordedSources.size, + }) + } added.push({ documentId: document.id, filename: document.filename, diff --git a/apps/sim/lib/knowledge/documents/service.ts b/apps/sim/lib/knowledge/documents/service.ts index 1db5a4b1a42..91b68123ad9 100644 --- a/apps/sim/lib/knowledge/documents/service.ts +++ b/apps/sim/lib/knowledge/documents/service.ts @@ -71,6 +71,7 @@ import { EXACT_EMPTY_DURABLE_SECRET_PROVENANCE, mergeDurableSecretProvenance, } from '@/lib/execution/durable-secret-provenance' +import { reportDurableSecretProvenanceUnrecorded } from '@/lib/execution/durable-secret-provenance-telemetry' import { knowledgeAccessCondition, knowledgeMetadataCandidateAccessCondition, @@ -662,15 +663,22 @@ const KNOWLEDGE_DOCUMENT_TAG_FIELDS = new Set([ 'boolean3', ]) +function countUnrecordedFileSources( + provenances: ReadonlyMap +): number { + let recordCount = 0 + for (const provenance of provenances.values()) { + if (provenance.status === 'unrecorded') recordCount += 1 + } + return recordCount +} + function durableSecretProvenanceFromWorkspaceFile( provenance: WorkspaceFileSecretProvenance, binding: FileMetadataRecord ): DurableSecretProvenance { - /** - * `unrecorded` is a more specific `unknown`, and this boundary has not opted into the workspace - * file surface's policy, so it keeps refusing exactly as it did. - */ - if (provenance.status !== 'exact') return { status: 'unknown' } + if (provenance.status === 'unknown') return provenance + if (provenance.status === 'unrecorded') return EXACT_EMPTY_DURABLE_SECRET_PROVENANCE return { status: 'exact', entries: provenance.entries.map((entry) => ({ @@ -2535,7 +2543,7 @@ export async function createDocumentRecords( const resolvedDocuments = await resolveServerKnownDocumentSizes(documents) const totalBytes = resolvedDocuments.reduce((sum, docData) => sum + (docData.fileSize || 0), 0) const admission = await resolveDocumentStorageAdmission(knowledgeBaseId, totalBytes) - const { returnData, storageNotification } = await db.transaction(async (tx) => { + const { returnData, storageNotification, unrecordedCount } = await db.transaction(async (tx) => { let storageNotification: DocumentStorageNotification | null = null await tx.execute(sql`SELECT 1 FROM knowledge_base WHERE id = ${knowledgeBaseId} FOR UPDATE`) @@ -2730,9 +2738,20 @@ export async function createDocumentRecords( .where(eq(knowledgeBase.id, knowledgeBaseId)) } - return { returnData, storageNotification } + return { + returnData, + storageNotification, + unrecordedCount: countUnrecordedFileSources(boundFileProvenanceById), + } }) + if (unrecordedCount > 0) { + reportDurableSecretProvenanceUnrecorded({ + surface: 'workspace-file', + workspaceId: admission.workspaceId, + recordCount: unrecordedCount, + }) + } if (storageNotification) { void maybeNotifyStorageLimitForBillingContext( storageNotification.context, @@ -3204,7 +3223,7 @@ export async function createSingleDocument( ...processedTags, } - const storageNotification = await db.transaction(async (tx) => { + const { storageNotification, unrecordedCount } = await db.transaction(async (tx) => { let storageNotification: DocumentStorageNotification | null = null if (options?.uploadedArtifact) { await claimKnowledgeUploadForAttachment(tx, options.uploadedArtifact.cleanupEventId) @@ -3336,9 +3355,19 @@ export async function createSingleDocument( }) } - return storageNotification + return { + storageNotification, + unrecordedCount: countUnrecordedFileSources(boundFileProvenanceById), + } }) + if (unrecordedCount > 0) { + reportDurableSecretProvenanceUnrecorded({ + surface: 'workspace-file', + workspaceId: admission.workspaceId, + recordCount: unrecordedCount, + }) + } if (storageNotification) { void maybeNotifyStorageLimitForBillingContext( storageNotification.context, diff --git a/apps/sim/lib/mothership/agent-cli/file-provenance.test.ts b/apps/sim/lib/mothership/agent-cli/file-provenance.test.ts index bb5d217bdbf..509d3239e9b 100644 --- a/apps/sim/lib/mothership/agent-cli/file-provenance.test.ts +++ b/apps/sim/lib/mothership/agent-cli/file-provenance.test.ts @@ -43,6 +43,7 @@ vi.mock('@/lib/execution/remote-sandbox/session-file-snapshot', () => ({ vi.mock('@/lib/mothership/chat/delegation', () => ({ mintDelegationToken: async () => 'fixture' })) import { withWorkspaceInvocationScope } from '@/lib/core/application/workspace-invocation-scope' +import { OrchestrationError } from '@/lib/core/orchestration/types' vi.mock('@/lib/mothership/application/workspace-target', () => mothershipWorkspaceTargetMock) vi.mock('@/lib/mothership/agent-cli/scoped-transport', () => ({ @@ -319,6 +320,38 @@ describe('file provenance at the actual CLI and model-result boundary', () => { ).not.toContain(content) }) + it('redacts generated document contents in a download renderer failure', async () => { + resetDbChainMock() + queueTableRows(workspaceFiles, [ + { + fileContentUpdatedAt: revision, + secretProvenanceVersion: 1, + provenanceContentUpdatedAt: revision, + status: 'exact', + entries: [ + { + name: 'FILE_SECRET', + encryptedValue: 'fixture-ciphertext', + sourceUserId: 'reader', + sourceWorkspaceId: 'workspace', + }, + ], + }, + ]) + mocks.file.mockResolvedValue({ ...file, name: 'private.pdf', type: 'text/x-pdflibjs' }) + mocks.render.mockRejectedValue( + new OrchestrationError('conflict', `Document could not be generated: ${content}`) + ) + const trace = registry() + const response = await readTransport(trace)(fileRequest()) + expect(response.status).toBe(409) + const output = await response.text() + expect(output).toContain(content) + const projected = inspectToolResultForCopilot({ success: false, output }, trace, 'sim_cli') + expect(JSON.stringify(projected.result)).not.toContain(content) + expect(JSON.stringify(projected.result)).toContain('{{FILE_SECRET}}') + }) + it.each(['sequential', 'parallel'])( 'carries file classification into each %s grep invocation', async (mode) => { diff --git a/apps/sim/lib/mothership/agent-cli/file-read-transport.ts b/apps/sim/lib/mothership/agent-cli/file-read-transport.ts index 85588ce18a3..cfccc827b0e 100644 --- a/apps/sim/lib/mothership/agent-cli/file-read-transport.ts +++ b/apps/sim/lib/mothership/agent-cli/file-read-transport.ts @@ -19,12 +19,14 @@ import { } from '@/lib/mothership/auth/application-delegation' import { importWorkspaceFileSnapshotProvenance, + mergeWorkspaceFileSecretProvenance, type WorkspaceFileSecretProvenance, } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' import { v2FileErrorPolicies } from '@/lib/workspace-files/api' import { presentWorkspaceFileText } from '@/lib/workspace-files/api/text-presenter' import { WORKSPACE_FILES_DELEGATION_AUDIENCE } from '@/lib/workspace-files/application/authorization' import { downloadWorkspaceFileStream } from '@/lib/workspace-files/application/download-workspace-file' +import { observeWorkspaceFileDelivery } from '@/lib/workspace-files/application/file-delivery-observer' import { readWorkspaceFileText } from '@/lib/workspace-files/application/read-workspace-file-text' import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' @@ -60,6 +62,40 @@ export function createFileReadTransport(context: { actorUserId: context.userId, })) if (!imported) context.registry?.markIncomplete('workspace-file-provenance-unknown') + return imported + } + + const forward: typeof fetch = async (input, init) => { + const url = new URL(input instanceof Request ? input.url : input) + const method = (init?.method ?? (input instanceof Request ? input.method : 'GET')).toUpperCase() + if (!context.invocation || method !== 'GET' || url.origin !== base.origin) { + return (context.transport ?? fetch)(input, init) + } + let provenance: WorkspaceFileSecretProvenance | undefined + const response = await observeWorkspaceFileDelivery( + async (source) => { + provenance = mergeWorkspaceFileSecretProvenance( + provenance ?? { status: 'exact', entries: [] }, + source ?? { status: 'unknown' } + ) + }, + () => (context.transport ?? fetch)(input, init) + ) + if (!provenance) return response + try { + if (!context.registry || !(await observe(context.invocation.workspaceId, provenance))) { + await response.body?.cancel().catch(() => {}) + return Response.json( + { error: { message: 'File read provenance is unavailable. Retry the read.' } }, + { status: 503 } + ) + } + if (response.body) context.trackDownload?.(response.body, provenance) + return response + } catch (error) { + await response.body?.cancel().catch(() => {}) + throw error + } } return async (input, init) => { @@ -71,10 +107,10 @@ export function createFileReadTransport(context: { !url.pathname.startsWith(prefix) || collectionPaths.has(url.pathname) ) { - return (context.transport ?? fetch)(input, init) + return forward(input, init) } const match = /^([^/]+)(\/text)?$/.exec(url.pathname.slice(prefix.length)) - if (!match) return (context.transport ?? fetch)(input, init) + if (!match) return forward(input, init) const request = new NextRequest(new Request(input, init)) let stream: ReadableStream | undefined try { @@ -106,20 +142,25 @@ export function createFileReadTransport(context: { ) if (!parsed.success) return parsed.response request.signal.throwIfAborted() - const result = await readWorkspaceFileText.execute({ - principal, - input: { - workspaceId: parsed.data.query.workspaceId, - reference: parsed.data.params.fileId, - ...(context.chatId !== undefined ? { chatId: context.chatId } : {}), - maxBytes: parsed.data.query.maxBytes, - offset: parsed.data.query.offset, - limit: parsed.data.query.limit, - includeSecretProvenance: true, - allowPlainText: true, + const result = await observeWorkspaceFileDelivery( + async (source) => { + await observe(parsed.data.query.workspaceId, source) }, - }) - await observe(result.file.workspaceId, result.secretProvenance) + () => + readWorkspaceFileText.execute({ + principal, + input: { + workspaceId: parsed.data.query.workspaceId, + reference: parsed.data.params.fileId, + ...(context.chatId !== undefined ? { chatId: context.chatId } : {}), + maxBytes: parsed.data.query.maxBytes, + offset: parsed.data.query.offset, + limit: parsed.data.query.limit, + includeSecretProvenance: true, + allowPlainText: true, + }, + }) + ) request.signal.throwIfAborted() return Response.json(presentWorkspaceFileText(result)) } @@ -131,16 +172,21 @@ export function createFileReadTransport(context: { ) if (!parsed.success) return parsed.response request.signal.throwIfAborted() - const result = await downloadWorkspaceFileStream.execute({ - principal, - input: { - fileId: parsed.data.params.fileId, - assertedWorkspaceId: parsed.data.query.workspaceId, - includeSecretProvenance: true, + const result = await observeWorkspaceFileDelivery( + async (source) => { + await observe(parsed.data.query.workspaceId, source) }, - }) + () => + downloadWorkspaceFileStream.execute({ + principal, + input: { + fileId: parsed.data.params.fileId, + assertedWorkspaceId: parsed.data.query.workspaceId, + includeSecretProvenance: true, + }, + }) + ) stream = result.stream - await observe(result.file.workspaceId, result.secretProvenance) request.signal.throwIfAborted() const response = new Response( stream.pipeThrough(new TransformStream(), { signal: request.signal }), diff --git a/apps/sim/lib/mothership/agent-cli/run-cli.ts b/apps/sim/lib/mothership/agent-cli/run-cli.ts index a3ed837a645..d44a161f33b 100644 --- a/apps/sim/lib/mothership/agent-cli/run-cli.ts +++ b/apps/sim/lib/mothership/agent-cli/run-cli.ts @@ -77,6 +77,7 @@ export async function readCliInputFile( const registry = new ResolvedSecretTraceRegistry([]) if (!(await importDurableSecretProvenance(registry, provenance))) throw new Error('CLI input withheld because workbench secret provenance is unavailable') + if (registry.getActiveMatches().length === 0) return buffer const text = buffer.toString('utf8') const projection = projectResolvedSecretModelContent(text, registry) if (!projection.safe || isBinarySandboxPath(path) || projection.value !== text) diff --git a/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.integration.ts b/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.integration.ts new file mode 100644 index 00000000000..733c078aada --- /dev/null +++ b/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.integration.ts @@ -0,0 +1,510 @@ +import { createHash } from 'node:crypto' +import * as audit from '@sim/audit' +import { AuditAction } from '@sim/audit' +import { db } from '@sim/db' +import { auditLog, organization, user, workspace } from '@sim/db/schema' +import { readTestRedisUrl } from '@sim/db/testing/test-infrastructure' +import { flushMacrotask } from '@sim/testing/helpers/async' +import { redisConfigMock, redisConfigMockFns } from '@sim/testing/mocks/redis-config.mock' +import { generateShortId } from '@sim/utils/id' +import { and, eq } from 'drizzle-orm' +import Redis from 'ioredis' +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' +import { buildOrgScopeCondition } from '@/lib/audit-logs/query' +import { initializeSessionFileProvenance } from '@/lib/execution/remote-sandbox/session-file-provenance' +import { createWorkbenchFileProvenance } from '@/lib/mothership/agent-cli/workbench-file-provenance' +import type { WorkspaceFileSecretProvenance } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' +import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' + +vi.mock('@/lib/core/config/redis', () => redisConfigMock) + +const redisUrl = readTestRedisUrl() +const redis = redisUrl + ? new Redis(redisUrl, { + lazyConnect: true, + retryStrategy: () => null, + maxRetriesPerRequest: 1, + connectTimeout: 1000, + }) + : undefined +redis?.on('error', () => {}) +const receiptPrefixes = new Set() +const historyKeys = new Set() +const fixtureWorkspaceId = generateShortId() +const fixtureUserId = generateShortId() +let scope = { workspaceId: fixtureWorkspaceId, userId: fixtureUserId, sessionKey: 'chat' } +let submittedAudits = 0 +let restoreAuditObservation: (() => void) | undefined +const machine = { providerId: 'e2b', sandboxId: 'physical-machine' } as const +const bytes = new Uint8Array([255, 254, 0, 1, 90, 13, 10]) +const secret: WorkspaceFileSecretProvenance = { + status: 'exact', + entries: [{ encryptedValue: 'ciphertext', sourceUserId: 'reader' }], +} +const safe: WorkspaceFileSecretProvenance = { status: 'exact', entries: [] } +const body = (content = bytes) => + new ReadableStream({ + start(controller) { + controller.enqueue(content) + controller.close() + }, + }) +const consume = async (stream: ReadableStream) => + new Uint8Array(await new Response(stream).arrayBuffer()) + +async function download(provenance: WorkspaceFileSecretProvenance = secret) { + const invocation = createWorkbenchFileProvenance(scope) + const stream = body() + invocation.trackDownload(stream, provenance) + expect(await consume(invocation.observeDownload(machine, stream))).toEqual(bytes) +} + +function receiptPrefix(currentScope: { + workspaceId?: string + organizationId?: string + userId: string + sessionKey: string +}) { + const namespace = createHash('sha256') + .update( + JSON.stringify([ + currentScope.organizationId + ? { organizationId: currentScope.organizationId } + : currentScope.workspaceId, + currentScope.userId, + currentScope.sessionKey, + machine.providerId, + machine.sandboxId, + ]) + ) + .digest('hex') + const prefix = `mothership:file-source:${currentScope.organizationId ? 'v2' : 'v1'}:${namespace}` + receiptPrefixes.add(prefix) + return prefix +} + +/** Waits for real fire-and-forget audit writes before inspecting or removing fixture rows. */ +async function settledWorkspaceAudits() { + const read = () => db.select().from(auditLog).where(eq(auditLog.workspaceId, fixtureWorkspaceId)) + await vi.waitFor(async () => expect(await read()).toHaveLength(submittedAudits)) + return read() +} + +beforeAll(async () => { + if (!redis) return + await redis.connect() + await db.insert(user).values({ + id: fixtureUserId, + name: 'Receipt test actor', + email: `${fixtureUserId}@fixture.test`, + emailVerified: true, + createdAt: new Date(), + updatedAt: new Date(), + }) + await db.insert(workspace).values({ + id: fixtureWorkspaceId, + name: 'Receipt test workspace', + ownerId: fixtureUserId, + billedAccountUserId: fixtureUserId, + }) + const recordAudit = audit.recordAudit + const observation = vi.spyOn(audit, 'recordAudit').mockImplementation((entry) => { + if (entry.workspaceId === fixtureWorkspaceId) submittedAudits += 1 + recordAudit(entry) + }) + restoreAuditObservation = () => observation.mockRestore() +}) +beforeEach(() => { + redisConfigMockFns.mockGetRedisClient.mockReturnValue(redis ?? null) + scope = { ...scope, sessionKey: `chat-${generateShortId(16)}` } + receiptPrefix(scope) + historyKeys.add( + `mothership:workbench-provenance:v2:${createHash('sha256') + .update(JSON.stringify([scope.sessionKey, machine.providerId, machine.sandboxId])) + .digest('hex')}` + ) +}) +afterEach(async () => { + if (!redis) return + await settledWorkspaceAudits() + await db.delete(auditLog).where(eq(auditLog.workspaceId, fixtureWorkspaceId)) + submittedAudits = 0 +}) +afterAll(async () => { + if (!redis) return + try { + for (const prefix of receiptPrefixes) { + let cursor = '0' + do { + const [next, keys] = await redis.scan(cursor, 'MATCH', `${prefix}:*`, 'COUNT', 100) + cursor = next + if (keys.length > 0) await redis.del(...keys) + } while (cursor !== '0') + } + if (historyKeys.size > 0) await redis.del(...historyKeys) + await db.delete(auditLog).where(eq(auditLog.workspaceId, fixtureWorkspaceId)) + await db.delete(workspace).where(eq(workspace.id, fixtureWorkspaceId)) + await db.delete(user).where(eq(user.id, fixtureUserId)) + } finally { + restoreAuditObservation?.() + await Promise.all([redis.quit(), db.$client.end()]) + } +}) + +describe.skipIf(!redisUrl)('trusted workbench byte receipts with real Redis', () => { + it.each(['complete', 'pending sibling', 'unknown'] as const)( + 'persists settled stdout evidence from a %s registry', + async (state) => { + const registry = new ResolvedSecretTraceRegistry([], scope) + if (state === 'unknown') registry.markIncomplete('workspace-file-provenance-unknown') + const settle = state === 'pending sibling' ? registry.beginPendingActivation() : undefined + const invocation = createWorkbenchFileProvenance({ + ...scope, + resolvedSecretTraceRegistry: registry, + }) + const value = 'saved command output' + try { + const observe = await invocation.observeOutput(value) + await consume(observe(machine, new Blob([value]).stream())) + const next = createWorkbenchFileProvenance(scope) + await consume(next.observeUpload(machine, new Blob([value]).stream())) + expect(next.uploadProvenance()).toEqual(state === 'unknown' ? { status: 'unknown' } : safe) + } finally { + settle?.() + } + } + ) + it('keeps a large safe transfer streamed and readable across invocations', async () => { + const size = 6 * 1024 * 1024 + 17 + const chunk = new Uint8Array(64 * 1024).fill(255) + const source = () => { + let remaining = size + return new ReadableStream({ + pull(controller) { + if (!remaining) { + controller.close() + return + } + const length = Math.min(remaining, chunk.length) + remaining -= length + controller.enqueue(chunk.subarray(0, length)) + }, + }) + } + const count = async (stream: ReadableStream) => { + let received = 0 + await stream.pipeTo( + new WritableStream({ + write(value) { + received += value.byteLength + }, + }) + ) + expect(received).toBe(size) + } + const first = createWorkbenchFileProvenance(scope) + const stream = source() + first.trackDownload(stream, safe) + await count(first.observeDownload(machine, stream)) + const next = createWorkbenchFileProvenance(scope) + await count(next.observeUpload(machine, source())) + expect(next.uploadProvenance()).toEqual(safe) + }) + it.each([secret, safe, { status: 'unknown' } as const])( + 'survives a fresh invocation: %j', + async (source) => { + await download(source) + const next = createWorkbenchFileProvenance(scope) + expect(() => next.uploadProvenance()).toThrow('has not finished') + expect(await consume(next.observeUpload(machine, body()))).toEqual(bytes) + expect(next.uploadProvenance()).toEqual(source) + } + ) + + it('classifies changed bytes as unknown without poisoning an unchanged safe copy', async () => { + await download(safe) + const changed = createWorkbenchFileProvenance(scope) + await consume(changed.observeUpload(machine, body(new Uint8Array([0])))) + expect(changed.uploadProvenance()).toEqual({ status: 'unknown' }) + const unchanged = createWorkbenchFileProvenance(scope) + await consume(unchanged.observeUpload(machine, body())) + expect(unchanged.uploadProvenance()).toEqual(safe) + }) + + it.each(['workspaceId', 'userId', 'sessionKey', 'sandboxId', 'providerId'] as const)( + 'isolates receipts by %s when machine history is unavailable', + async (field) => { + await download() + await redis!.del(...historyKeys) + const next = createWorkbenchFileProvenance({ + ...scope, + ...(['workspaceId', 'userId', 'sessionKey'].includes(field) ? { [field]: 'other' } : {}), + }) + await consume( + next.observeUpload( + { + ...machine, + ...(field === 'sandboxId' ? { sandboxId: 'replacement' } : {}), + ...(field === 'providerId' ? { providerId: 'daytona' as const } : {}), + }, + body() + ) + ) + expect(next.uploadProvenance()).toEqual({ status: 'unknown' }) + } + ) + + it.each([ + ['unrecorded then empty', { status: 'unrecorded' }, safe, safe], + ['empty then unrecorded', safe, { status: 'unrecorded' }, safe], + ['unrecorded then known', { status: 'unrecorded' }, secret, secret], + ['known then unrecorded', secret, { status: 'unrecorded' }, secret], + [ + 'unrecorded then unknown', + { status: 'unrecorded' }, + { status: 'unknown' }, + { status: 'unknown' }, + ], + [ + 'unknown then unrecorded', + { status: 'unknown' }, + { status: 'unrecorded' }, + { status: 'unknown' }, + ], + ] as const)( + 'preserves recorded evidence for identical bytes: %s', + async (_label, first, second, expected) => { + await download(first) + await download(second) + const next = createWorkbenchFileProvenance(scope) + expect(await consume(next.observeUpload(machine, body()))).toEqual(bytes) + expect(next.uploadProvenance()).toEqual(expected) + } + ) + + it('concurrent conflicting receipts cannot overwrite unknown with safe', async () => { + await Promise.all([download(safe), download({ status: 'unknown' }), download(secret)]) + await download(safe) + const next = createWorkbenchFileProvenance(scope) + await consume(next.observeUpload(machine, body())) + expect(next.uploadProvenance()).toEqual({ status: 'unknown' }) + }) + + it('parallel independent streams keep their own classification', async () => { + const invocation = createWorkbenchFileProvenance(scope) + const first = body() + const second = body(new Uint8Array([1])) + invocation.trackDownload(first, secret) + invocation.trackDownload(second, safe) + await Promise.all([ + consume(invocation.observeDownload(machine, first)), + consume(invocation.observeDownload(machine, second)), + ]) + const next = createWorkbenchFileProvenance(scope) + await consume(next.observeUpload(machine, body())) + expect(next.uploadProvenance()).toEqual(secret) + const other = createWorkbenchFileProvenance(scope) + await consume(other.observeUpload(machine, body(new Uint8Array([1])))) + expect(other.uploadProvenance()).toEqual(safe) + }) + + it('cannot complete from a partially consumed or failed stream', async () => { + const invocation = createWorkbenchFileProvenance(scope) + const failing = new ReadableStream({ + start(controller) { + controller.enqueue(bytes) + }, + pull(controller) { + controller.error(new Error('transfer failed')) + }, + }) + await expect(consume(invocation.observeUpload(machine, failing))).rejects.toThrow( + 'transfer failed' + ) + expect(() => invocation.uploadProvenance()).toThrow('has not finished') + }) + + it('Stop makes even a fully read source unavailable for completion', async () => { + await download() + const controller = new AbortController() + const invocation = createWorkbenchFileProvenance({ ...scope, signal: controller.signal }) + await consume(invocation.observeUpload(machine, body())) + controller.abort(new Error('stopped')) + expect(() => invocation.uploadProvenance()).toThrow('stopped') + }) + + it('missing history and present invalid receipts never certify bytes safe', async () => { + const next = createWorkbenchFileProvenance(scope) + await consume(next.observeUpload(machine, body())) + expect(next.uploadProvenance()).toEqual({ status: 'unknown' }) + await initializeSessionFileProvenance(scope.sessionKey, machine) + const digest = createHash('sha256').update(bytes).digest('hex') + const key = `${receiptPrefix(scope)}:${digest}` + for (const value of [ + 'broken JSON', + JSON.stringify({ + version: 1, + workspaceId: scope.workspaceId, + provenance: { status: 'unknown' }, + }), + ]) { + await redis!.set(key, value, 'EX', 60) + const broken = createWorkbenchFileProvenance(scope) + await consume(broken.observeUpload(machine, body())) + expect(broken.uploadProvenance()).toEqual({ status: 'unknown' }) + } + redisConfigMockFns.mockGetRedisClient.mockReturnValue(null) + const unavailable = createWorkbenchFileProvenance(scope) + await expect(consume(unavailable.observeUpload(machine, body()))).rejects.toThrow( + 'storage is unavailable' + ) + expect(() => unavailable.uploadProvenance()).toThrow('has not finished') + }) + + it.each([ + 'complete', + 'source failure', + 'cancel', + 'abort', + 'storage failure', + 'abort during persistence', + ] as const)('audits unrecorded downloads only after acceptance: %s', async (outcome) => { + const stop = new AbortController() + let producer!: ReadableStreamDefaultController + const stream = new ReadableStream({ + start(controller) { + producer = controller + controller.enqueue(bytes) + }, + }) + const invocation = createWorkbenchFileProvenance({ ...scope, signal: stop.signal }) + invocation.trackDownload(stream, { status: 'unrecorded' }) + const reader = invocation.observeDownload(machine, stream).getReader() + try { + expect(await reader.read()).toEqual({ done: false, value: bytes }) + expect(await settledWorkspaceAudits()).toEqual([]) + if (outcome === 'cancel') { + await reader.cancel() + } else if (outcome === 'source failure') { + producer.error(new Error('transfer failed')) + await expect(reader.read()).rejects.toThrow('transfer failed') + } else if (outcome === 'abort') { + stop.abort(new Error('stopped')) + await expect(reader.read()).rejects.toThrow('stopped') + } else { + if (outcome === 'storage failure') { + redisConfigMockFns.mockGetRedisClient.mockReturnValueOnce(null) + } else if (outcome === 'abort during persistence') { + redisConfigMockFns.mockGetRedisClient.mockImplementationOnce(() => { + stop.abort(new Error('stopped during persistence')) + return redis! + }) + } + producer.close() + if (outcome === 'complete') { + expect(await reader.read()).toEqual({ done: true, value: undefined }) + } else { + await expect(reader.read()).rejects.toThrow( + outcome === 'storage failure' ? 'storage is unavailable' : 'stopped during persistence' + ) + } + } + if (outcome === 'abort during persistence') { + /** Cancellation rejects the reader before the already-issued receipt write settles. */ + await redis!.ping() + await flushMacrotask() + } + const entries = await settledWorkspaceAudits() + expect(entries).toHaveLength(outcome === 'complete' ? 1 : 0) + if (outcome === 'complete') { + expect(entries[0]).toMatchObject({ + action: AuditAction.SECRET_PROVENANCE_UNRECORDED, + actorId: fixtureUserId, + workspaceId: fixtureWorkspaceId, + }) + const next = createWorkbenchFileProvenance(scope) + await consume(next.observeUpload(machine, body())) + expect(next.uploadProvenance()).toEqual({ status: 'unrecorded' }) + } + } finally { + await reader.cancel().catch(() => {}) + reader.releaseLock() + } + }) + + it('records accepted organization bytes in the owning organization audit scope', async () => { + const organizationId = generateShortId() + const userId = generateShortId() + await db + .insert(organization) + .values({ id: organizationId, name: 'Receipt audit fixture', slug: organizationId }) + await db.insert(user).values({ + id: userId, + name: 'Receipt audit actor', + email: `${userId}@fixture.test`, + emailVerified: true, + createdAt: new Date(), + updatedAt: new Date(), + }) + try { + const orgScope = { organizationId, userId, sessionKey: scope.sessionKey } + receiptPrefix(orgScope) + const invocation = createWorkbenchFileProvenance(orgScope) + const stream = body() + invocation.trackDownload(stream, { status: 'unrecorded' }) + expect(await consume(invocation.observeDownload(machine, stream))).toEqual(bytes) + const query = (owner: string) => + db + .select({ + workspaceId: auditLog.workspaceId, + actorId: auditLog.actorId, + metadata: auditLog.metadata, + }) + .from(auditLog) + .where( + and( + eq(auditLog.action, AuditAction.SECRET_PROVENANCE_UNRECORDED), + buildOrgScopeCondition({ + organizationId: owner, + orgWorkspaceIds: [], + orgMemberIds: [userId], + includeDeparted: false, + }) + ) + ) + await expect.poll(() => query(organizationId)).toHaveLength(1) + expect(await query(organizationId)).toEqual([ + { + workspaceId: null, + actorId: userId, + metadata: { surface: 'workspace-file', organizationId }, + }, + ]) + expect(await query(generateShortId())).toEqual([]) + } finally { + await db.delete(auditLog).where(eq(auditLog.actorId, userId)) + await db.delete(user).where(eq(user.id, userId)) + await db.delete(organization).where(eq(organization.id, organizationId)) + } + }) + + it('shares an org chat receipt across explicit targets but never across orgs or chats', async () => { + const orgScope = { ...scope, organizationId: 'org', workspaceId: 'a' } + receiptPrefix(orgScope) + const first = createWorkbenchFileProvenance(orgScope) + const stream = body() + first.trackDownload(stream, safe) + await consume(first.observeDownload(machine, stream)) + const second = createWorkbenchFileProvenance({ ...orgScope, workspaceId: 'b' }) + await consume(second.observeUpload(machine, body())) + expect(second.uploadProvenance()).toEqual(safe) + for (const other of [ + { ...orgScope, organizationId: 'other' }, + { ...orgScope, sessionKey: 'other' }, + ]) { + const outsider = createWorkbenchFileProvenance(other) + await consume(outsider.observeUpload(machine, body())) + expect(outsider.uploadProvenance()).toEqual({ status: 'unknown' }) + } + }) +}) diff --git a/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.test.ts b/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.test.ts deleted file mode 100644 index cb2ad044ab5..00000000000 --- a/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.test.ts +++ /dev/null @@ -1,269 +0,0 @@ -import { createHash } from 'node:crypto' -import { redisConfigMockFns } from '@sim/testing/mocks/redis-config.mock' -import { generateShortId } from '@sim/utils/id' -import Redis from 'ioredis' -import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest' -import { createWorkbenchFileProvenance } from '@/lib/mothership/agent-cli/workbench-file-provenance' -import type { WorkspaceFileSecretProvenance } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' -import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' - -/** The same assertions can exercise the actual Lua against an explicitly supplied disposable socket. */ -const redis = process.env.MSHIP_TEST_REDIS_SOCKET - ? new Redis({ - path: process.env.MSHIP_TEST_REDIS_SOCKET, - lazyConnect: true, - retryStrategy: () => null, - maxRetriesPerRequest: 1, - connectTimeout: 1000, - }) - : undefined -redis?.on('error', () => {}) -const recorded = new Map() -const memory = { - eval: vi.fn( - async (_script: string, _count: number, key: string, value: string, unknown: string) => { - const previous = recorded.get(key) - recorded.set(key, previous && previous !== value ? unknown : value) - return 1 - } - ), - get: vi.fn(async (key: string) => recorded.get(key) ?? null), -} -let storage: typeof memory | Redis | null = redis ?? memory -redisConfigMockFns.mockGetRedisClient.mockImplementation(() => storage) -let scope = { workspaceId: 'workspace', userId: 'reader', sessionKey: 'chat' } -const machine = { providerId: 'e2b', sandboxId: 'physical-machine' } as const -const bytes = new Uint8Array([255, 254, 0, 1, 90, 13, 10]) -const secret: WorkspaceFileSecretProvenance = { - status: 'exact', - entries: [{ encryptedValue: 'ciphertext', sourceUserId: 'reader' }], -} -const safe: WorkspaceFileSecretProvenance = { status: 'exact', entries: [] } -const body = (content = bytes) => - new ReadableStream({ - start(controller) { - controller.enqueue(content) - controller.close() - }, - }) -const consume = async (stream: ReadableStream) => - new Uint8Array(await new Response(stream).arrayBuffer()) - -async function download(provenance: WorkspaceFileSecretProvenance = secret) { - const invocation = createWorkbenchFileProvenance(scope) - const stream = body() - invocation.trackDownload(stream, provenance) - expect(await consume(invocation.observeDownload(machine, stream))).toEqual(bytes) -} - -beforeEach(() => { - recorded.clear() - storage = redis ?? memory - scope = { ...scope, sessionKey: `chat-${generateShortId(16)}` } -}) -afterAll(async () => { - await redis?.quit() -}) - -describe('trusted workbench byte receipts', () => { - it.each(['complete', 'pending sibling', 'unknown'] as const)( - 'persists settled stdout evidence from a %s registry', - async (state) => { - const registry = new ResolvedSecretTraceRegistry([], scope) - if (state === 'unknown') registry.markIncomplete('workspace-file-provenance-unknown') - const settle = state === 'pending sibling' ? registry.beginPendingActivation() : undefined - const invocation = createWorkbenchFileProvenance({ - ...scope, - resolvedSecretTraceRegistry: registry, - }) - const value = 'saved command output' - try { - const observe = await invocation.observeOutput(value) - await consume(observe(machine, new Blob([value]).stream())) - const next = createWorkbenchFileProvenance(scope) - await consume(next.observeUpload(machine, new Blob([value]).stream())) - expect(next.uploadProvenance()).toEqual(state === 'unknown' ? { status: 'unknown' } : safe) - } finally { - settle?.() - } - } - ) - it('keeps a large safe transfer streamed and readable across invocations', async () => { - const size = 6 * 1024 * 1024 + 17 - const chunk = new Uint8Array(64 * 1024).fill(255) - const source = () => { - let remaining = size - return new ReadableStream({ - pull(controller) { - if (!remaining) { - controller.close() - return - } - const length = Math.min(remaining, chunk.length) - remaining -= length - controller.enqueue(chunk.subarray(0, length)) - }, - }) - } - const count = async (stream: ReadableStream) => { - let received = 0 - await stream.pipeTo( - new WritableStream({ - write(value) { - received += value.byteLength - }, - }) - ) - expect(received).toBe(size) - } - const first = createWorkbenchFileProvenance(scope) - const stream = source() - first.trackDownload(stream, safe) - await count(first.observeDownload(machine, stream)) - const next = createWorkbenchFileProvenance(scope) - await count(next.observeUpload(machine, source())) - expect(next.uploadProvenance()).toEqual(safe) - }) - it.each([secret, safe, { status: 'unknown' } as const])( - 'survives a fresh invocation: %j', - async (source) => { - await download(source) - const next = createWorkbenchFileProvenance(scope) - expect(() => next.uploadProvenance()).toThrow('has not finished') - expect(await consume(next.observeUpload(machine, body()))).toEqual(bytes) - expect(next.uploadProvenance()).toEqual(source) - } - ) - - it('classifies changed bytes as unknown without poisoning an unchanged safe copy', async () => { - await download(safe) - const changed = createWorkbenchFileProvenance(scope) - await consume(changed.observeUpload(machine, body(new Uint8Array([0])))) - expect(changed.uploadProvenance()).toEqual({ status: 'unknown' }) - const unchanged = createWorkbenchFileProvenance(scope) - await consume(unchanged.observeUpload(machine, body())) - expect(unchanged.uploadProvenance()).toEqual(safe) - }) - - it.each(['workspaceId', 'userId', 'sessionKey', 'sandboxId', 'providerId'] as const)( - 'isolates evidence by %s', - async (field) => { - await download() - const next = createWorkbenchFileProvenance({ - ...scope, - ...(['workspaceId', 'userId', 'sessionKey'].includes(field) ? { [field]: 'other' } : {}), - }) - await consume( - next.observeUpload( - { - ...machine, - ...(field === 'sandboxId' ? { sandboxId: 'replacement' } : {}), - ...(field === 'providerId' ? { providerId: 'daytona' as const } : {}), - }, - body() - ) - ) - expect(next.uploadProvenance()).toEqual({ status: 'unknown' }) - } - ) - - it('concurrent conflicting receipts cannot overwrite unknown with safe', async () => { - await Promise.all([download(safe), download({ status: 'unknown' }), download(secret)]) - await download(safe) - const next = createWorkbenchFileProvenance(scope) - await consume(next.observeUpload(machine, body())) - expect(next.uploadProvenance()).toEqual({ status: 'unknown' }) - }) - - it('parallel independent streams keep their own classification', async () => { - const invocation = createWorkbenchFileProvenance(scope) - const first = body() - const second = body(new Uint8Array([1])) - invocation.trackDownload(first, secret) - invocation.trackDownload(second, safe) - await Promise.all([ - consume(invocation.observeDownload(machine, first)), - consume(invocation.observeDownload(machine, second)), - ]) - const next = createWorkbenchFileProvenance(scope) - await consume(next.observeUpload(machine, body())) - expect(next.uploadProvenance()).toEqual(secret) - const other = createWorkbenchFileProvenance(scope) - await consume(other.observeUpload(machine, body(new Uint8Array([1])))) - expect(other.uploadProvenance()).toEqual(safe) - }) - - it('cannot complete from a partially consumed or failed stream', async () => { - const invocation = createWorkbenchFileProvenance(scope) - const failing = new ReadableStream({ - start(controller) { - controller.enqueue(bytes) - }, - pull(controller) { - controller.error(new Error('transfer failed')) - }, - }) - await expect(consume(invocation.observeUpload(machine, failing))).rejects.toThrow( - 'transfer failed' - ) - expect(() => invocation.uploadProvenance()).toThrow('has not finished') - }) - - it('Stop makes even a fully read source unavailable for completion', async () => { - await download() - const controller = new AbortController() - const invocation = createWorkbenchFileProvenance({ ...scope, signal: controller.signal }) - await consume(invocation.observeUpload(machine, body())) - controller.abort(new Error('stopped')) - expect(() => invocation.uploadProvenance()).toThrow('stopped') - }) - - it('missing, malformed and unavailable storage never certifies bytes safe', async () => { - const next = createWorkbenchFileProvenance(scope) - await consume(next.observeUpload(machine, body())) - expect(next.uploadProvenance()).toEqual({ status: 'unknown' }) - const namespace = createHash('sha256') - .update( - JSON.stringify([ - scope.workspaceId, - scope.userId, - scope.sessionKey, - machine.providerId, - machine.sandboxId, - ]) - ) - .digest('hex') - const digest = createHash('sha256').update(bytes).digest('hex') - const key = `mothership:file-source:v1:${namespace}:${digest}` - if (redis) await redis.set(key, 'broken JSON', 'EX', 60) - else recorded.set(key, 'broken JSON') - const broken = createWorkbenchFileProvenance(scope) - await consume(broken.observeUpload(machine, body())) - expect(broken.uploadProvenance()).toEqual({ status: 'unknown' }) - storage = null - const unavailable = createWorkbenchFileProvenance(scope) - await expect(consume(unavailable.observeUpload(machine, body()))).rejects.toThrow( - 'storage is unavailable' - ) - expect(() => unavailable.uploadProvenance()).toThrow('has not finished') - }) -}) - -it('shares an org chat receipt across explicit targets but never across orgs or chats', async () => { - const orgScope = { ...scope, organizationId: 'org', workspaceId: 'a' } - const first = createWorkbenchFileProvenance(orgScope) - const stream = body() - first.trackDownload(stream, safe) - await consume(first.observeDownload(machine, stream)) - const second = createWorkbenchFileProvenance({ ...orgScope, workspaceId: 'b' }) - await consume(second.observeUpload(machine, body())) - expect(second.uploadProvenance()).toEqual(safe) - for (const other of [ - { ...orgScope, organizationId: 'other' }, - { ...orgScope, sessionKey: 'other' }, - ]) { - const outsider = createWorkbenchFileProvenance(other) - await consume(outsider.observeUpload(machine, body())) - expect(outsider.uploadProvenance()).toEqual({ status: 'unknown' }) - } -}) diff --git a/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.ts b/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.ts index 8bef5617c0e..40f1f7ba4c0 100644 --- a/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.ts +++ b/apps/sim/lib/mothership/agent-cli/workbench-file-provenance.ts @@ -1,10 +1,15 @@ import { createHash } from 'node:crypto' import { getRedisClient } from '@/lib/core/config/redis' +import { createDurableSecretProvenanceRegistry } from '@/lib/execution/durable-secret-provenance' +import { reportDurableSecretProvenanceUnrecorded } from '@/lib/execution/durable-secret-provenance-telemetry' import type { SessionFileIdentity, SessionFileObserver, } from '@/lib/execution/remote-sandbox/session-file-observer' -import { recordSessionFileInput } from '@/lib/execution/remote-sandbox/session-file-provenance' +import { + readSessionSecretProvenance, + recordSessionFileInput, +} from '@/lib/execution/remote-sandbox/session-file-provenance' import { createWorkspaceFileSecretProvenanceFromRegistry, type WorkspaceFileSecretProvenance, @@ -17,12 +22,14 @@ import { } from '@/lib/uploads/upload-session/workspace-file-provenance' import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' -/** Missing/expired evidence is unknown. Redis is never a source of exact-empty by default. */ +/** Hash receipts retain source classifications for unchanged bytes across invocations. */ const RECEIPT_SECONDS = 24 * 60 * 60 const RECORD_RECEIPT = ` local previous = redis.call('GET', KEYS[1]) local value = ARGV[1] -if previous and previous ~= value then value = ARGV[2] end +if previous and previous ~= value and previous ~= ARGV[4] then + if value == ARGV[4] then value = previous else value = ARGV[2] end +end redis.call('SET', KEYS[1], value, 'EX', ARGV[3]) return 1 ` @@ -64,6 +71,7 @@ export function createWorkbenchFileProvenance(scope: WorkbenchFileScope) { : bindWorkspaceFileUploadProvenance(scope.workspaceId!, provenance) ) const unknown = encoded({ status: 'unknown' }) + const unrecorded = encoded({ status: 'unrecorded' }) const record = (provenance: WorkspaceFileSecretProvenance): SessionFileObserver => @@ -80,24 +88,62 @@ export function createWorkbenchFileProvenance(scope: WorkbenchFileScope) { key(machine, digest), encoded(provenance), unknown, - RECEIPT_SECONDS + RECEIPT_SECONDS, + unrecorded ) + scope.signal?.throwIfAborted() + if (provenance.status === 'unrecorded') { + reportDurableSecretProvenanceUnrecorded({ + surface: 'workspace-file', + workspaceId: scope.workspaceId, + organizationId: scope.organizationId, + actorUserId: scope.userId, + }) + } }, () => recordSessionFileInput( scope.sessionKey, machine, - provenance.status === 'exact' ? provenance : { status: 'unknown' } + provenance.status === 'unrecorded' ? true : provenance ) ) const observeDownload: SessionFileObserver = (machine, stream) => - record(downloads.get(stream) ?? { status: 'unknown' })(machine, stream) + record(downloads.get(stream) ?? { status: 'unrecorded' })(machine, stream) const observeUpload: SessionFileObserver = (machine, stream) => { uploaded = undefined return hashStream(stream, scope.signal, async (digest) => { const redis = getRedisClient() if (!redis) throw new Error('Workbench file classification storage is unavailable') const value = await redis.get(key(machine, digest)) + if (value === null) { + /** Newly generated bytes inherit settled machine history; an existing bad receipt never falls back. */ + const history = await readSessionSecretProvenance(scope.sessionKey, machine) + if (history.status === 'unknown') { + uploaded = history + return + } + try { + const registry = await createDurableSecretProvenanceRegistry(history, scope) + uploaded = !registry?.getActiveMatches().length + ? { status: 'exact', entries: [] } + : parseWorkspaceFileSecretProvenance({ + status: 'exact', + entries: history.entries.map((entry) => + entry.sourceUserId + ? entry + : { + encryptedValue: entry.encryptedValue, + sourceUserId: scope.userId, + ...(scope.workspaceId ? { sourceWorkspaceId: scope.workspaceId } : {}), + } + ), + }) + } catch { + uploaded = { status: 'unknown' } + } + return + } let binding: unknown try { binding = value ? JSON.parse(value) : undefined diff --git a/apps/sim/lib/mothership/chat/application/read-sandbox-file.test.ts b/apps/sim/lib/mothership/chat/application/read-sandbox-file.test.ts index 07e5875fc24..98bd6590cee 100644 --- a/apps/sim/lib/mothership/chat/application/read-sandbox-file.test.ts +++ b/apps/sim/lib/mothership/chat/application/read-sandbox-file.test.ts @@ -11,7 +11,6 @@ import { createCopilotChatFilePrincipal } from '@/lib/mothership/auth/file-deleg const hoisted = vi.hoisted(() => ({ context: vi.fn(), snapshot: vi.fn(), - clean: vi.fn(), receipt: vi.fn(), dispose: vi.fn(), })) @@ -19,9 +18,6 @@ vi.mock('@sim/platform-authz/workspace', () => workspaceAuthzMock) vi.mock('@/lib/mothership/chat/application/context', () => ({ resolveOwnedChatContext: hoisted.context, })) -vi.mock('@/lib/execution/remote-sandbox/session-file-provenance', () => ({ - isSessionFileProvenanceClean: hoisted.clean, -})) vi.mock('@/lib/execution/remote-sandbox/session-file-snapshot', () => ({ openSessionFileSnapshot: hoisted.snapshot, })) @@ -55,8 +51,7 @@ beforeEach(() => { allowPersonalApiKeys: true, billedAccountUserId: 'owner', }) - mocks.clean.mockResolvedValue(true) - mocks.receipt.mockReturnValue({ status: 'unknown' }) + mocks.receipt.mockReturnValue({ status: 'exact', entries: [] }) mocks.snapshot.mockImplementation(async (_key, _path, _signal, observer) => ({ size: 3, dispose: hoisted.dispose, @@ -87,7 +82,6 @@ describe('owned scratch snapshot reads', () => { expect.any(Function), { allowedRoots: ['/home/user', '/tmp'], maxBytes: 25 * 1024 * 1024 } ) - expect(mocks.clean).toHaveBeenCalledWith('mothership-chat:chat', machine) expect(mocks.dispose).toHaveBeenCalledOnce() }) it.each(['permission', 'foreign-workspace', 'foreign-chat'])( @@ -103,24 +97,12 @@ describe('owned scratch snapshot reads', () => { it.each([ { status: 'unknown' }, { status: 'exact', entries: [{ encryptedValue: 'secret', sourceUserId: 'u' }] }, - ])( - 'refuses unknown or secret prior machine inputs even without a current registry: %j', - async (receipt) => { - mocks.clean.mockResolvedValue(false) - mocks.receipt.mockReturnValue(receipt) - await expect(readChatSandboxFile.execute({ principal, input })).rejects.toThrow( - 'no verified secret-free provenance' - ) - expect(mocks.dispose).toHaveBeenCalledOnce() - } - ) - it('permits an unchanged exact-empty digest receipt after other machine inputs became unknown', async () => { - mocks.clean.mockResolvedValue(false) - mocks.receipt.mockReturnValue({ status: 'exact', entries: [] }) - expect((await readChatSandboxFile.execute({ principal, input })).buffer).toEqual( - Buffer.from('png') + ])('refuses explicit unknown or secret-bearing opaque file provenance: %j', async (receipt) => { + mocks.receipt.mockReturnValue(receipt) + await expect(readChatSandboxFile.execute({ principal, input })).rejects.toThrow( + 'no verified secret-free provenance' ) - expect(mocks.clean).not.toHaveBeenCalled() + expect(mocks.dispose).toHaveBeenCalledOnce() }) it('enforces size before streaming and still disposes the snapshot', async () => { await expect( @@ -176,17 +158,6 @@ describe('scoped chat file delegation', () => { ) }) -it('never lets a broad clean marker override an exact secret-bearing file receipt', async () => { - mocks.clean.mockResolvedValue(true) - mocks.receipt.mockReturnValue({ - status: 'exact', - entries: [{ encryptedValue: 'ciphertext', sourceUserId: 'u' }], - }) - await expect(readChatSandboxFile.execute({ principal, input })).rejects.toThrow( - 'no verified secret-free provenance' - ) -}) - it('accepts the verified personal API principal through the same canonical chat authorization', async () => { const personal = createPersonalApiKeyPrincipal({ userId: 'u', keyId: 'verified-key' }) expect((await readChatSandboxFile.execute({ principal: personal, input })).name).toBe('image.png') @@ -201,3 +172,10 @@ it('rejects a workspace API key before canonical lookup or sandbox access', asyn expect(mocks.context).not.toHaveBeenCalled() expect(mocks.snapshot).not.toHaveBeenCalled() }) + +it('permits a source explicitly classified as unrecorded', async () => { + mocks.receipt.mockReturnValue({ status: 'unrecorded' }) + expect((await readChatSandboxFile.execute({ principal, input })).buffer).toEqual( + Buffer.from('png') + ) +}) diff --git a/apps/sim/lib/mothership/chat/application/read-sandbox-file.ts b/apps/sim/lib/mothership/chat/application/read-sandbox-file.ts index dbd4fd04351..cd5b9fe8543 100644 --- a/apps/sim/lib/mothership/chat/application/read-sandbox-file.ts +++ b/apps/sim/lib/mothership/chat/application/read-sandbox-file.ts @@ -4,8 +4,7 @@ import { defineWorkspaceOperation } from '@/lib/core/application' import { defineOrganizationOperation } from '@/lib/core/application/organization-operation' import { OrchestrationError } from '@/lib/core/orchestration/types' import { readStreamToBufferWithLimit } from '@/lib/core/utils/stream-limits' -import type { SessionFileIdentity } from '@/lib/execution/remote-sandbox/session-file-observer' -import { isSessionFileProvenanceClean } from '@/lib/execution/remote-sandbox/session-file-provenance' +import { reportDurableSecretProvenanceUnrecorded } from '@/lib/execution/durable-secret-provenance-telemetry' import { openSessionFileSnapshot } from '@/lib/execution/remote-sandbox/session-file-snapshot' import { createWorkbenchFileProvenance } from '@/lib/mothership/agent-cli/workbench-file-provenance' import { defineAuthorizedChatUseCase } from '@/lib/mothership/chat/application/authorized-chat-use-case' @@ -93,15 +92,11 @@ export const readChatSandboxFile = defineAuthorizedChatUseCase({ sessionKey, signal: input.signal, }) - let machine: SessionFileIdentity | undefined const snapshot = await openSessionFileSnapshot( sessionKey, path, input.signal, - (identity, stream) => { - machine = identity - return provenance.observeUpload(identity, stream) - }, + provenance.observeUpload, { allowedRoots: ['/home/user', '/tmp'], maxBytes } ) try { @@ -116,10 +111,9 @@ export const readChatSandboxFile = defineAuthorizedChatUseCase({ }) input.signal?.throwIfAborted() const receipt = provenance.uploadProvenance() - const exactEmpty = receipt.status === 'exact' && receipt.entries.length === 0 if ( - (receipt.status === 'exact' && receipt.entries.length > 0) || - (!exactEmpty && (!machine || !(await isSessionFileProvenanceClean(sessionKey, machine)))) + receipt.status === 'unknown' || + (receipt.status === 'exact' && receipt.entries.length > 0) ) { throw new OrchestrationError( 'forbidden', @@ -127,6 +121,14 @@ export const readChatSandboxFile = defineAuthorizedChatUseCase({ ) } input.signal?.throwIfAborted() + if (receipt.status === 'unrecorded') { + reportDurableSecretProvenanceUnrecorded({ + surface: 'workspace-file', + workspaceId: context.workspaceId, + organizationId: context.organizationId, + actorUserId: context.userId, + }) + } return { buffer, path, name: posix.basename(path) } } finally { await snapshot.dispose() diff --git a/apps/sim/lib/mothership/tools/handlers/function-execute-file-mounts.test.ts b/apps/sim/lib/mothership/tools/handlers/function-execute-file-mounts.test.ts index 5d8d3e2f210..dee936c49ad 100644 --- a/apps/sim/lib/mothership/tools/handlers/function-execute-file-mounts.test.ts +++ b/apps/sim/lib/mothership/tools/handlers/function-execute-file-mounts.test.ts @@ -430,24 +430,63 @@ describe('organization-owned upload mounts', () => { expect(mocks.list).not.toHaveBeenCalled() expect(registry.isComplete()).toBe(true) }) - it('rejects unknown provenance without certifying derived scratch', async () => { - const registry = new ResolvedSecretTraceRegistry([], { userId: 'reader' }) - registry.markIncomplete('mounted-file-provenance-unavailable') - await expect( - resolveInputFiles( + it.each(['missing', 'absence', 'matching-secret'] as const)( + 'accepts intentionally uploaded bytes with $0 source provenance without activating them', + async (source) => { + const parent = new ResolvedSecretTraceRegistry([], { userId: 'reader' }) + if (source === 'absence') parent.markIncomplete('source-provenance-incomplete') + if (source === 'matching-secret') { + mocks.decrypt.mockResolvedValue({ decrypted: 'public input' }) + await parent.importProvenance( + { + version: 1, + complete: true, + scope: { userId: 'reader' }, + entries: [{ name: 'TOKEN', encryptedValue: 'encrypted:public input' }], + }, + { trusted: true } + ) + } + const registry = new ResolvedSecretTraceRegistry([], { userId: 'reader' }) + const result = await resolveInputFiles( { ...context, workspaceId: undefined, organizationId: 'org', - resolvedSecretTraceRegistry: registry, + ...(source === 'missing' ? {} : { resolvedSecretTraceRegistry: parent }), }, - [{ path: 'uploads/upload' }], + [{ path: 'uploads/upload', sandboxPath: '/tmp/input.txt' }], [], [], registry ) - ).rejects.toThrow('provenance') - expect(registry.isComplete()).toBe(false) + expect(result).toEqual([ + { + path: '/tmp/input.txt', + content: Buffer.from('public input').toString('base64'), + encoding: 'base64', + }, + ]) + expect(registry.exportCheckpointProvenance()).toMatchObject({ complete: true, entries: [] }) + } + ) + it('does not repair an existing protection fault when an intentional upload is mounted', async () => { + const registry = new ResolvedSecretTraceRegistry([], { userId: 'reader' }) + registry.markIncomplete('entry-decrypt-failed') + const result = await resolveInputFiles( + { + ...context, + workspaceId: undefined, + organizationId: 'org', + resolvedSecretTraceRegistry: registry, + }, + [{ path: 'uploads/upload' }], + [], + [], + registry + ) + expect(result).toHaveLength(1) + expect(registry.isPermanentlyIncomplete()).toBe(true) }) it('requires a target for workspace files and never treats missing upload authority as a workspace fallback', async () => { const registry = new ResolvedSecretTraceRegistry([], { userId: 'reader' }) diff --git a/apps/sim/lib/mothership/tools/handlers/function-execute.ts b/apps/sim/lib/mothership/tools/handlers/function-execute.ts index 50e6281582a..2574ac2ec08 100644 --- a/apps/sim/lib/mothership/tools/handlers/function-execute.ts +++ b/apps/sim/lib/mothership/tools/handlers/function-execute.ts @@ -5,7 +5,6 @@ import { hasWorkspaceSandboxAccess } from '@/lib/billing/core/subscription' import { OrchestrationError } from '@/lib/core/orchestration/types' import { durableSecretProvenanceFromEnvelope, - importDurableSecretProvenance, mergeDurableSecretProvenance, } from '@/lib/execution/durable-secret-provenance' import type { PrivateSecretProvenanceBundleV1 } from '@/lib/execution/model-input-provenance' @@ -54,10 +53,7 @@ import { parseChatUploadReference, type WorkspaceFileRecord, } from '@/lib/uploads/contexts/workspace/workspace-file-manager' -import { - createWorkspaceFileSecretProvenanceFromRegistry, - importWorkspaceFileSnapshotProvenance, -} from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' +import { importWorkspaceFileSnapshotProvenance } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' import { WORKSPACE_FILES_DELEGATION_AUDIENCE } from '@/lib/workspace-files/application/authorization' import { listAllWorkspaceFiles } from '@/lib/workspace-files/application/list-workspace-files' import { fileOperations } from '@/lib/workspace-files/application/operations' @@ -275,23 +271,6 @@ export async function resolveInputFiles( byteCount += file.buffer.length if (byteCount > MAX_INLINE_MOUNT_TOTAL_BYTES) throw new Error('Chat attachment mounts exceed the buffered mount limit') - if (!resolvedSecretTraceRegistry) throw new Error('Chat attachment provenance is unavailable') - const classification = await createWorkspaceFileSecretProvenanceFromRegistry( - context.resolvedSecretTraceRegistry, - { text: file.buffer.toString('utf8'), base64: file.buffer.toString('base64') }, - { userId: context.userId, ...(workspaceId ? { workspaceId } : {}) } - ) - if (!classification.safe || classification.provenance.status === 'unknown') - throw new Error('Chat attachment provenance is unavailable') - if (classification.provenance.status === 'exact') { - if ( - !(await importDurableSecretProvenance( - resolvedSecretTraceRegistry, - classification.provenance - )) - ) - throw new Error('Chat attachment provenance could not be imported') - } else resolvedSecretTraceRegistry.markIncomplete('mounted-file-provenance-unavailable') attachments.push({ path: refField(fileRef, 'sandboxPath') ?? diff --git a/apps/sim/lib/mothership/tools/handlers/workbench-confidentiality.live.test.ts b/apps/sim/lib/mothership/tools/handlers/workbench-confidentiality.live.test.ts index 83042c5f896..94cf2baa8c8 100644 --- a/apps/sim/lib/mothership/tools/handlers/workbench-confidentiality.live.test.ts +++ b/apps/sim/lib/mothership/tools/handlers/workbench-confidentiality.live.test.ts @@ -1,6 +1,8 @@ import { execFile } from 'node:child_process' +import { createReadStream } from 'node:fs' import { mkdir, mkdtemp, readdir, readFile, rm, stat, writeFile } from 'node:fs/promises' import { dirname, join } from 'node:path' +import { Readable } from 'node:stream' import { promisify } from 'node:util' import { createDelegatedPrincipal } from '@sim/testing/factories/principal.factory' import { createDeferred } from '@sim/testing/helpers/deferred' @@ -61,6 +63,7 @@ vi.mock('@/lib/mothership/vfs/resource-writer', () => ({ import { functionExecuteBodySchema } from '@/lib/api/contracts' import * as inProcessTransport from '@/lib/api/server/routes/in-process-transport' import { encryptSecret } from '@/lib/core/security/encryption' +import { importDurableSecretProvenance } from '@/lib/execution/durable-secret-provenance' import { PRIVATE_TOOL_METADATA_REQUEST_HEADER, RESOLVED_SECRET_NAMES_FIELD, @@ -75,6 +78,7 @@ import { import type { SandboxHandle } from '@/lib/execution/remote-sandbox/types' import { executeFunctionRequest } from '@/lib/function-execution/execute-request' import { readCliInputFile } from '@/lib/mothership/agent-cli/run-cli' +import { createWorkbenchFileProvenance } from '@/lib/mothership/agent-cli/workbench-file-provenance' import { inspectToolResultForCopilot } from '@/lib/mothership/request/tools/resolved-secret-result' import type { ToolExecutionContext } from '@/lib/mothership/tool-executor/types' import { executeFunctionExecute } from '@/lib/mothership/tools/handlers/function-execute' @@ -88,6 +92,7 @@ import { buildMothershipSandboxSession } from '@/lib/mothership/tools/sandbox-se import { chatSandboxSessionKey } from '@/lib/mothership/tools/sandbox-session-key' import { reportTableRowDelivery } from '@/lib/table/application/row-delivery-observer' import { reportWorkspaceFileDelivery } from '@/lib/workspace-files/application/file-delivery-observer' +import { projectResolvedSecretModelContent } from '@/executor/utils/resolved-secret-content-projection' import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' import { buildFunctionExecuteBody, functionExecuteTool } from '@/tools/function/execute' import type { CodeExecutionInput } from '@/tools/function/types' @@ -353,6 +358,170 @@ async function sandboxApi(path: string, handler: () => Promise, method }) } +describe('generated workbench file provenance', () => { + const identity = () => ({ providerId: 'e2b' as const, sandboxId: machine.sandboxId }) + const observer = () => + createWorkbenchFileProvenance({ + ...scope, + sessionKey: chatSandboxSessionKey(chatId), + }) + const consume = async (stream: ReadableStream) => + Buffer.from(await new Response(stream).arrayBuffer()) + + it.each(['receipt.json', 'screenshot.png'])( + 'admits fresh %s bytes from a clean workbench', + async (path) => { + const bytes = path.endsWith('.json') + ? Buffer.from('{"done":true}') + : Buffer.from([137, 80, 78, 71, 13, 10, 26, 10]) + const result = await run( + `printf '%s' '${bytes.toString('base64')}' | base64 --decode > ${path}` + ) + expect(result.raw.success).toBe(true) + expect(result.projected.safe).toBe(true) + const file = observer() + const stream = Readable.toWeb( + createReadStream(workerPath(`/home/user/${path}`)) + ) as ReadableStream + expect(await consume(file.observeUpload(identity(), stream))).toEqual(bytes) + expect(file.uploadProvenance()).toEqual({ status: 'exact', entries: [] }) + } + ) + + it.each(['unrecorded', 'unregistered'] as const)( + 'allows %s downloaded bytes without losing earlier named secret protection', + async (classification) => { + await run('printf "%s" "$TOKEN" > saved.txt', ['TOKEN']) + const file = observer() + const stream = new Blob(['ordinary user input']).stream() + if (classification === 'unrecorded') file.trackDownload(stream, { status: 'unrecorded' }) + await machine.writeFile( + '/home/user/input.txt', + await consume(file.observeDownload(identity(), stream)) + ) + const result = await run('cat input.txt saved.txt') + expect(result.projected.safe).toBe(true) + expect(result.projected.result).toMatchObject({ + success: true, + output: { stdout: 'ordinary user input{{TOKEN}}' }, + }) + } + ) + + it('keeps an expected but missing download classification unavailable', async () => { + const file = observer() + const stream = new Blob(['ordinary bytes']).stream() + file.trackDownload(stream, undefined) + await consume(file.observeDownload(identity(), stream)) + expect(await readSessionSecretProvenance(chatSandboxSessionKey(chatId), identity())).toEqual({ + status: 'unknown', + }) + }) + + it('does not clear an earlier protection fault when unrecorded bytes arrive', async () => { + await recordSessionFileInput(chatSandboxSessionKey(chatId), identity(), false) + const file = observer() + await consume(file.observeDownload(identity(), new Blob(['ordinary bytes']).stream())) + expect(await readSessionSecretProvenance(chatSandboxSessionKey(chatId), identity())).toEqual({ + status: 'unknown', + }) + }) + + it.each(['same', 'foreign', 'anonymous'] as const)( + 'carries %s-scope machine secrets into new file redaction', + async (source) => { + await recordSessionFileInput(chatSandboxSessionKey(chatId), identity(), { + status: 'exact', + entries: [ + { + name: 'TOKEN', + encryptedValue: catalog[0].encryptedValue, + ...(source === 'anonymous' + ? {} + : { + sourceUserId: source === 'same' ? scope.userId : 'another-user', + sourceWorkspaceId: scope.workspaceId, + }), + }, + ], + }) + const file = observer() + await consume(file.observeUpload(identity(), new Blob([canary]).stream())) + const provenance = file.uploadProvenance() + expect(provenance.status).toBe('exact') + if (provenance.status !== 'exact') throw new Error('Expected verified provenance') + const registry = new ResolvedSecretTraceRegistry([], scope) + expect(await importDurableSecretProvenance(registry, provenance)).toBe(true) + expect(projectResolvedSecretModelContent(canary, registry)).toMatchObject({ + safe: true, + value: source === 'same' ? '{{TOKEN}}' : '[REDACTED_SECRET]', + }) + } + ) + + it('never treats failed history decryption as an empty protected-secret set', async () => { + await recordSessionFileInput(chatSandboxSessionKey(chatId), identity(), { + status: 'exact', + entries: [{ encryptedValue: 'invalid-ciphertext', sourceUserId: scope.userId }], + }) + const file = observer() + await consume(file.observeUpload(identity(), new Blob(['ordinary file']).stream())) + expect(file.uploadProvenance()).toEqual({ status: 'unknown' }) + }) + + it('applies the shared short-value exemption to generated opaque bytes and CLI input', async () => { + const short = '1234567' + const encryptedValue = (await encryptSecret(short)).encrypted + await recordSessionFileInput(chatSandboxSessionKey(chatId), identity(), { + status: 'exact', + entries: [ + { + name: 'SHORT_TOKEN', + encryptedValue, + sourceUserId: scope.userId, + sourceWorkspaceId: scope.workspaceId, + }, + ], + }) + const bytes = Buffer.from([137, 80, 78, 71, 13, 10, 26, 10]) + await machine.writeFile('/home/user/image.png', bytes) + expect(await readCliInputFile(chatSandboxSessionKey(chatId), 'image.png')).toEqual(bytes) + const file = observer() + await consume(file.observeUpload(identity(), new Blob([bytes]).stream())) + expect(file.uploadProvenance()).toEqual({ status: 'exact', entries: [] }) + }) + + it('includes secrets admitted while upload bytes are still streaming', async () => { + const began = createDeferred() + const finish = createDeferred() + const file = observer() + const stream = new ReadableStream({ + async start(controller) { + controller.enqueue(Buffer.from(canary)) + began.resolve() + await finish.promise + controller.close() + }, + }) + const consumed = consume(file.observeUpload(identity(), stream)) + await began.promise + await recordSessionFileInput(chatSandboxSessionKey(chatId), identity(), { + status: 'exact', + entries: [ + { + name: 'TOKEN', + encryptedValue: catalog[0].encryptedValue, + sourceUserId: scope.userId, + sourceWorkspaceId: scope.workspaceId, + }, + ], + }) + finish.resolve() + await consumed + expect(file.uploadProvenance()).toMatchObject({ status: 'exact', entries: [{ name: 'TOKEN' }] }) + }) +}) + describe('sandbox API provenance admission', () => { it('keeps ordinary API mutations usable for later code output and generated CLI input', async () => { const response = await sandboxApi( @@ -425,6 +594,21 @@ describe('sandbox API provenance admission', () => { } ) + it('accepts explicitly unrecorded file delivery and retains earlier secret history', async () => { + await run('printf "%s" "$TOKEN" > saved.txt', ['TOKEN']) + const response = await sandboxApi('/api/v2/files/fixture/download', async () => { + await reportWorkspaceFileDelivery({ status: 'unrecorded' }) + return new Response('ordinary user input') + }) + await machine.writeFile('/home/user/input.txt', await response.text()) + const result = await run('cat input.txt saved.txt') + expect(result.projected.safe).toBe(true) + expect(result.projected.result).toMatchObject({ + success: true, + output: { stdout: 'ordinary user input{{TOKEN}}' }, + }) + }) + it('preserves mutation completion and withholds its body when provenance storage fails', async () => { const mutationPath = join(root, 'mutation.json') const response = await sandboxApi( diff --git a/apps/sim/lib/mothership/tools/sandbox-resource-transport.test.ts b/apps/sim/lib/mothership/tools/sandbox-resource-transport.test.ts index 39ea3ef05b8..12781b3eb2b 100644 --- a/apps/sim/lib/mothership/tools/sandbox-resource-transport.test.ts +++ b/apps/sim/lib/mothership/tools/sandbox-resource-transport.test.ts @@ -7,7 +7,6 @@ import { beforeEach, describe, expect, it, vi } from 'vitest' import { isCopilotRequest } from '@/lib/api/server/routes/copilot-request' import { assertWorkspaceInvocationScope } from '@/lib/core/application/workspace-invocation-scope' import { OrchestrationError } from '@/lib/core/orchestration/types' -import { reportWorkspaceFileDelivery } from '@/lib/workspace-files/application/file-delivery-observer' const { readScope, recordEffects, fetcher, routeMatcher, recordInput, mint } = vi.hoisted(() => ({ readScope: vi.fn(), @@ -266,33 +265,6 @@ it('does not poison public scratch after authenticated static catalog discovery' expect(recordInput).not.toHaveBeenCalled() }) -it.each([ - { status: 'exact', entries: [] }, - { status: 'unknown' }, - { status: 'exact', entries: [{ encryptedValue: 'ciphertext', sourceUserId: 'user' }] }, -] as const)( - 'uses typed file provenance from the real public handler without changing admission: %j', - async (provenance) => { - routeMatcher.mockReturnValue({ - params: { fileId: 'file' }, - load: async () => ({ GET: fetcher }), - }) - fetcher.mockImplementation(async (input: Request) => { - await reportWorkspaceFileDelivery({ - ...provenance, - ...('entries' in provenance ? { entries: [...provenance.entries] } : {}), - }) - return new Response('filebytes') - }) - const response = await proxySandboxResourceRequest(request('/api/v2/files/file'), token) - expect(await response.text()).toBe('filebytes') - expect(recordInput).toHaveBeenCalledWith( - 'mothership-chat:chat', - provenance.status === 'exact' ? provenance : false - ) - expect(fetcher).toHaveBeenCalledOnce() - } -) it('cannot deliver a successful unclassified file response when recording its unknown provenance fails', async () => { routeMatcher.mockReturnValue({ params: { fileId: 'file' }, load: async () => ({ GET: fetcher }) }) fetcher.mockResolvedValue(new Response('unsafe')) diff --git a/apps/sim/lib/mothership/tools/sandbox-resource-transport.ts b/apps/sim/lib/mothership/tools/sandbox-resource-transport.ts index e35d6ffaf9b..2f5f23ec218 100644 --- a/apps/sim/lib/mothership/tools/sandbox-resource-transport.ts +++ b/apps/sim/lib/mothership/tools/sandbox-resource-transport.ts @@ -13,6 +13,7 @@ import { EXACT_EMPTY_DURABLE_SECRET_PROVENANCE, mergeDurableSecretProvenance, } from '@/lib/execution/durable-secret-provenance' +import { reportDurableSecretProvenanceUnrecorded } from '@/lib/execution/durable-secret-provenance-telemetry' import { recordExistingSessionFileInput } from '@/lib/execution/remote-sandbox/session-file-provenance' import { createResourceEffectTransport } from '@/lib/mothership/agent-cli/resource-effects' import { resolveInvocationWorkspace } from '@/lib/mothership/application/workspace-target' @@ -146,7 +147,14 @@ async function proxyAuthorizedSandboxRequest( }, () => observeWorkspaceFileDelivery(async (provenance) => { - await recordInput(provenance?.status === 'exact' ? provenance : false) + if (provenance?.status === 'unrecorded') { + reportDurableSecretProvenanceUnrecorded({ + surface: 'workspace-file', + workspaceId: targetWorkspaceId, + actorUserId: scope.userId, + }) + } + await recordInput(provenance?.status === 'unrecorded' ? true : (provenance ?? false)) fileObserved = true }, dispatch) ) diff --git a/apps/sim/lib/uploads/archive.ts b/apps/sim/lib/uploads/archive.ts index cb449efa7b0..27675986a12 100644 --- a/apps/sim/lib/uploads/archive.ts +++ b/apps/sim/lib/uploads/archive.ts @@ -291,8 +291,8 @@ function throwInflateCapError(reason: 'entry' | 'total', entryName: string): nev * Filesystem-noise entries (`__MACOSX/`, `.DS_Store`, `Thumbs.db`) are extracted * verbatim unless `skipNoiseEntries` is set — the HTTP decompress route preserves * them; the agent-facing extract path drops them. Decompression is not byte-preserving, - * so only an exact-empty archive classification can remain exact on extracted files; - * every other classification becomes unknown without changing the extracted bytes. + * so known secret contributions become unknown on extracted files. Exact-empty and unrecorded + * classifications retain their existing input policy without changing the extracted bytes. */ export async function decompressArchiveBufferToWorkspaceFiles( buffer: Buffer, @@ -322,7 +322,8 @@ export async function decompressArchiveBufferToWorkspaceFiles( notifyWorkspaceChange = true, } = opts const extractedSecretProvenance: WorkspaceFileSecretProvenance = - secretProvenance.status === 'exact' && secretProvenance.entries.length === 0 + secretProvenance.status === 'unrecorded' || + (secretProvenance.status === 'exact' && secretProvenance.entries.length === 0) ? secretProvenance : { status: 'unknown' } diff --git a/apps/sim/lib/uploads/contexts/workspace/__integration__/file-versions.integration.ts b/apps/sim/lib/uploads/contexts/workspace/__integration__/file-versions.integration.ts index 5b839aad66c..e88fdc4cab6 100644 --- a/apps/sim/lib/uploads/contexts/workspace/__integration__/file-versions.integration.ts +++ b/apps/sim/lib/uploads/contexts/workspace/__integration__/file-versions.integration.ts @@ -24,10 +24,13 @@ vi.mock('@/lib/uploads/core/setup.server', () => ({ }, })) +import { encryptSecret } from '@/lib/core/security/encryption' import { createKnowledgeAclFixtureIds, seedKnowledgeAclFixture, } from '@/lib/knowledge/__integration__/seed-source-access-fixture' +import { createFileReadTransport } from '@/lib/mothership/agent-cli/file-read-transport' +import { runCli } from '@/lib/mothership/agent-cli/run-cli' import { deleteWorkspaceFileVersion, fetchWorkspaceFileBuffer, @@ -36,15 +39,25 @@ import { updateWorkspaceFileContent, uploadWorkspaceFile, } from '@/lib/uploads/contexts/workspace/workspace-file-manager' +import type { WorkspaceFileSecretProvenance } from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' import { WORKSPACE_FILE_STORAGE_CLEANUP_OUTBOX_EVENT } from '@/lib/uploads/contexts/workspace/workspace-file-storage-cleanup-outbox' import { getCurrentWorkspaceFileVersion, getWorkspaceFileVersion, + getWorkspaceFileVersionProvenance, queryWorkspaceFileVersions, releaseWorkspaceFileVersionsForPurgeInTx, } from '@/lib/uploads/contexts/workspace/workspace-file-versions' -import { revertWorkspaceFileVersion } from '@/lib/workspace-files/application/file-versions' +import { v2FileErrorPolicies } from '@/lib/workspace-files/api' +import { presentWorkspaceFileText } from '@/lib/workspace-files/api/text-presenter' +import { + downloadWorkspaceFileVersion, + readWorkspaceFileVersionText, + revertWorkspaceFileVersion, +} from '@/lib/workspace-files/application/file-versions' import { runCleanupFileVersions } from '@/background/cleanup-file-versions' +import { projectResolvedSecretModelContent } from '@/executor/utils/resolved-secret-content-projection' +import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' describe('workspace file version history in PostgreSQL', () => { const fixtures: ReturnType[] = [] @@ -63,7 +76,7 @@ describe('workspace file version history in PostgreSQL', () => { await Promise.all([db.$client.end(), dbFor('cleanup').$client.end()]) }) - async function seedFile(content: string) { + async function seedFile(content: string, extension = 'txt') { const ids = createKnowledgeAclFixtureIds() fixtures.push(ids) await seedKnowledgeAclFixture(ids) @@ -71,8 +84,8 @@ describe('workspace file version history in PostgreSQL', () => { ids.workspaceId, ids.aliceId, Buffer.from(content), - `notes-${generateId()}.txt`, - 'text/plain', + `notes-${generateId()}.${extension}`, + extension === 'json' ? 'application/json' : 'text/plain', { notifyWorkspaceChange: false } ) return { ...ids, fileId: uploaded.id, firstKey: uploaded.key } @@ -119,6 +132,323 @@ describe('workspace file version history in PostgreSQL', () => { expect(await versionRows(fixture.fileId)).toEqual([]) }) + it.each([ + ['named', 'download'], + ['anonymous', 'download'], + ['named', 'read'], + ['anonymous', 'read'], + ['empty', 'read'], + ['unrecorded', 'read'], + ] as const)( + 'preserves %s historical provenance through CLI versions %s', + async (kind, command) => { + const fixture = await seedFile('original') + const secret = 'SYNTHETIC_VERSION_SECRET_FOR_LOCAL_TEST' + const hasSecret = kind === 'named' || kind === 'anonymous' + const historicalContent = hasSecret ? secret : 'ordinary historical text' + const provenance: WorkspaceFileSecretProvenance = + kind === 'unrecorded' + ? { status: 'unrecorded' } + : { + status: 'exact', + entries: hasSecret + ? [ + { + encryptedValue: (await encryptSecret(secret)).encrypted, + sourceUserId: fixture.aliceId, + sourceWorkspaceId: fixture.workspaceId, + ...(kind === 'named' ? { name: 'FILE_TOKEN' } : {}), + }, + ] + : [], + } + for (const [content, secretProvenance] of [ + [historicalContent, provenance], + ['public replacement', { status: 'exact', entries: [] }], + ] as const) { + await updateWorkspaceFileContent( + fixture.workspaceId, + fixture.fileId, + fixture.aliceId, + Buffer.from(content), + undefined, + { + version: { source: 'api', authorUserId: fixture.aliceId }, + secretProvenancePolicy: { mode: 'replace', provenance: secretProvenance }, + } + ) + } + const registry = new ResolvedSecretTraceRegistry([], { + userId: fixture.aliceId, + workspaceId: fixture.workspaceId, + }) + const delivered: Array<{ + stream: ReadableStream + provenance: WorkspaceFileSecretProvenance | undefined + }> = [] + const transport = createFileReadTransport({ + endpoint: 'https://version-fixture.test', + userId: fixture.aliceId, + registry, + invocation: { + userId: fixture.aliceId, + workspaceId: fixture.workspaceId, + chatId: 'version-fixture', + }, + trackDownload: (stream, provenance) => { + delivered.push({ stream, provenance }) + }, + transport: async () => { + if (command === 'read') { + const result = await readWorkspaceFileVersionText.execute({ + principal: { kind: 'session', userId: fixture.aliceId, sessionId: generateId() }, + input: { + fileId: fixture.fileId, + assertedWorkspaceId: fixture.workspaceId, + version: 2, + }, + }) + return Response.json({ + data: { ...presentWorkspaceFileText(result).data, version: result.version.version }, + }) + } + const result = await downloadWorkspaceFileVersion.execute({ + principal: { kind: 'session', userId: fixture.aliceId, sessionId: generateId() }, + input: { fileId: fixture.fileId, assertedWorkspaceId: fixture.workspaceId, version: 2 }, + }) + return new Response(result.stream, { headers: { 'content-type': result.contentType } }) + }, + }) + const result = await runCli( + ['files', 'versions', command, fixture.fileId, '2'], + { + endpoint: 'https://version-fixture.test', + apiKey: 'fixture', + workspaceId: fixture.workspaceId, + transport, + }, + null + ) + expect(result.exitCode).toBe(0) + expect(result.stdout).toContain(historicalContent) + const projected = projectResolvedSecretModelContent(result.stdout, registry) + expect(projected.safe).toBe(true) + if (!projected.safe) throw new Error('Historical file output was withheld') + if (hasSecret) expect(projected.value).not.toContain(secret) + else expect(projected.value).toBe(result.stdout) + if (kind === 'named') expect(projected.value).toContain('{{FILE_TOKEN}}') + expect(delivered.map((delivery) => delivery.provenance)).toEqual([provenance]) + } + ) + + it.each(['legacy', 'nonempty', 'malformed'] as const)( + 'requires a valid empty legacy snapshot before delivering historical bytes: %s', + async (kind) => { + const content = 'Historical content whose protection must not be discarded' + const fixture = await seedFile(content) + await updateWorkspaceFileContent( + fixture.workspaceId, + fixture.fileId, + fixture.aliceId, + Buffer.from('current content'), + undefined, + { version: { source: 'api', authorUserId: fixture.aliceId } } + ) + await db + .update(workspaceFileVersion) + .set({ + secretProvenanceStatus: null, + secretProvenanceEntries: + kind === 'malformed' + ? sql`'{}'::jsonb` + : kind === 'nonempty' + ? [ + { + name: 'TOKEN', + encryptedValue: 'fixture-ciphertext', + sourceUserId: fixture.aliceId, + }, + ] + : [], + }) + .where( + and(eq(workspaceFileVersion.fileId, fixture.fileId), eq(workspaceFileVersion.version, 1)) + ) + const registry = new ResolvedSecretTraceRegistry([], { + userId: fixture.aliceId, + workspaceId: fixture.workspaceId, + }) + const transport = createFileReadTransport({ + endpoint: 'https://version-fixture.test', + userId: fixture.aliceId, + registry, + invocation: { + userId: fixture.aliceId, + workspaceId: fixture.workspaceId, + chatId: 'version-fixture', + }, + transport: async () => { + const result = await downloadWorkspaceFileVersion.execute({ + principal: { kind: 'session', userId: fixture.aliceId, sessionId: generateId() }, + input: { fileId: fixture.fileId, assertedWorkspaceId: fixture.workspaceId, version: 1 }, + }) + return new Response(result.stream, { headers: { 'content-type': result.contentType } }) + }, + }) + const result = await runCli( + ['files', 'versions', 'download', fixture.fileId, '1'], + { + endpoint: 'https://version-fixture.test', + apiKey: 'fixture', + workspaceId: fixture.workspaceId, + transport, + }, + null + ) + if (kind === 'legacy') { + expect(result.exitCode).toBe(0) + expect(result.stdout).toContain(content) + } else { + expect(result.exitCode).not.toBe(0) + expect(`${result.stdout}${result.stderr}`).not.toContain(content) + expect(`${result.stdout}${result.stderr}`).toContain('provenance is unavailable') + } + } + ) + + it.each(['current', 'historical'] as const)( + 'redacts %s file contents when the real parser includes them in a CLI error', + async (revision) => { + const fixture = await seedFile('{}', 'json') + const secret = 'ZZTOKEN99' + const provenance: WorkspaceFileSecretProvenance = { + status: 'exact', + entries: [ + { + encryptedValue: (await encryptSecret(secret)).encrypted, + sourceUserId: fixture.aliceId, + sourceWorkspaceId: fixture.workspaceId, + name: 'FILE_TOKEN', + }, + ], + } + await updateWorkspaceFileContent( + fixture.workspaceId, + fixture.fileId, + fixture.aliceId, + Buffer.from(secret), + 'application/json', + { + version: { source: 'api', authorUserId: fixture.aliceId }, + secretProvenancePolicy: { mode: 'replace', provenance }, + } + ) + if (revision === 'historical') { + await updateWorkspaceFileContent( + fixture.workspaceId, + fixture.fileId, + fixture.aliceId, + Buffer.from('{}'), + 'application/json', + { + version: { source: 'api', authorUserId: fixture.aliceId }, + secretProvenancePolicy: { + mode: 'replace', + provenance: { status: 'exact', entries: [] }, + }, + } + ) + } + const registry = new ResolvedSecretTraceRegistry([], { + userId: fixture.aliceId, + workspaceId: fixture.workspaceId, + }) + const transport = createFileReadTransport({ + endpoint: 'https://version-fixture.test', + userId: fixture.aliceId, + registry, + invocation: { + userId: fixture.aliceId, + workspaceId: fixture.workspaceId, + chatId: 'version-fixture', + }, + transport: async () => { + try { + const result = await readWorkspaceFileVersionText.execute({ + principal: { kind: 'session', userId: fixture.aliceId, sessionId: generateId() }, + input: { + fileId: fixture.fileId, + assertedWorkspaceId: fixture.workspaceId, + version: 2, + }, + }) + return Response.json({ + data: { ...presentWorkspaceFileText(result).data, version: result.version.version }, + }) + } catch (error) { + const response = v2FileErrorPolicies.concealResourceAuthorization.render(error) + if (!response) throw error + return response + } + }, + }) + const result = await runCli( + revision === 'current' + ? ['files', 'read', fixture.fileId] + : ['files', 'versions', 'read', fixture.fileId, '2'], + { + endpoint: 'https://version-fixture.test', + apiKey: 'fixture', + workspaceId: fixture.workspaceId, + transport, + }, + null + ) + expect(result.exitCode).not.toBe(0) + expect(result.stderr).toContain(secret) + const projected = projectResolvedSecretModelContent(result.stderr, registry) + expect(projected.safe).toBe(true) + if (!projected.safe) throw new Error('File parser error was withheld') + expect(projected.value).not.toContain(secret) + expect(projected.value).toContain('{{FILE_TOKEN}}') + } + ) + + it('does not bind a coalesced version classification to its previously captured storage key', async () => { + const fixture = await seedFile('original') + const write = { source: 'collab', authorUserId: fixture.aliceId } as const + const first = await updateWorkspaceFileContent( + fixture.workspaceId, + fixture.fileId, + fixture.aliceId, + Buffer.from('first draft'), + undefined, + { + version: write, + secretProvenancePolicy: { mode: 'replace', provenance: { status: 'unknown' } }, + } + ) + const captured = await getCurrentWorkspaceFileVersion(first) + const replacement = await updateWorkspaceFileContent( + fixture.workspaceId, + fixture.fileId, + fixture.aliceId, + Buffer.from('second draft'), + undefined, + { + version: write, + secretProvenancePolicy: { mode: 'replace', provenance: { status: 'exact', entries: [] } }, + } + ) + expect((await getCurrentWorkspaceFileVersion(replacement)).version).toBe(captured.version) + expect( + await getWorkspaceFileVersionProvenance(fixture.fileId, captured.version, captured.key) + ).toBeNull() + expect( + await getWorkspaceFileVersionProvenance(fixture.fileId, captured.version, replacement.key) + ).toEqual({ status: 'exact', entries: [] }) + }) + it('materializes version 1 on the first write and keeps the outgoing bytes readable', async () => { const fixture = await seedFile('original') diff --git a/apps/sim/lib/uploads/contexts/workspace/workspace-file-secret-provenance.test.ts b/apps/sim/lib/uploads/contexts/workspace/workspace-file-secret-provenance.test.ts index 4345e869751..72ddda46dec 100644 --- a/apps/sim/lib/uploads/contexts/workspace/workspace-file-secret-provenance.test.ts +++ b/apps/sim/lib/uploads/contexts/workspace/workspace-file-secret-provenance.test.ts @@ -13,6 +13,8 @@ vi.mock('@/lib/uploads/contexts/workspace/workspace-file-versions', () => ({ })) vi.mock('@/lib/execution/durable-secret-provenance-telemetry', () => ({ + reportDurableSecretProvenanceUnrecorded: vi.fn(), + reportDurableSecretProvenanceUnrecordedBatch: vi.fn(), reportDurableSecretProvenanceWrite: mockReportWrite, reportDurableSecretProvenanceRefusal: mockReportRefusal, })) @@ -379,6 +381,7 @@ describe('workspace file secret provenance', () => { * stored `unknown` above is dropped: a writer refused those bytes on purpose, which is a * different claim from nobody having recorded them, and no policy relaxes it. */ + { id: 'unrecorded-id', key: 'unrecorded-key' }, { id: 'pre-marker-sidecar-id', key: 'pre-marker-sidecar-key' }, { id: 'legacy-id', key: 'legacy-key' }, { id: 'inline-file' }, @@ -612,7 +615,7 @@ describe('workspace file secret provenance', () => { * Unrecorded says exactly what an untracked file says, and that one has always mounted. There is * nothing to import either way, so the mount proceeds and the workspace is told. */ - it('refuses to mount an unrecorded tracked file', async () => { + it('admits an unrecorded tracked file without importing secret entries', async () => { const registry = { importProvenance: vi.fn(), isPermanentlyIncomplete: vi.fn().mockReturnValue(false), @@ -633,7 +636,7 @@ describe('workspace file secret provenance', () => { identity: { fileId: 'file-1', key: 'file-key', context: 'workspace' }, registry, }) - ).resolves.toBe(false) + ).resolves.toBe(true) expect(registry.importProvenance).not.toHaveBeenCalled() }) @@ -795,8 +798,8 @@ describe('workspace file secret provenance', () => { entries: [{ name: 'TOKEN', encryptedValue: 'encrypted', sourceUserId: 'user-1' }], } const unrecorded = { status: 'unrecorded' as const } - expect(mergeWorkspaceFileSecretProvenance(known, unrecorded)).toEqual({ status: 'unknown' }) - expect(mergeWorkspaceFileSecretProvenance(unrecorded, known)).toEqual({ status: 'unknown' }) + expect(mergeWorkspaceFileSecretProvenance(known, unrecorded)).toEqual(known) + expect(mergeWorkspaceFileSecretProvenance(unrecorded, known)).toEqual(known) }) it('fails closed when persisted provenance is malformed', async () => { @@ -972,7 +975,7 @@ describe('workspace file secret provenance', () => { ) }) - it('refuses an unrecorded tracked file', async () => { + it('admits an unrecorded tracked file', async () => { queueTableRows(workspaceFiles, [ { id: 'unrecorded-id', @@ -987,12 +990,7 @@ describe('workspace file secret provenance', () => { }, ]) - await expect(isModelSafeWorkspaceFileKey('unrecorded-key')).resolves.toBe(false) - expect(mockReportRefusal).toHaveBeenCalledWith({ - surface: 'workspace-file', - cause: 'workspace-file-unrecorded-enforced', - workspaceId: undefined, - }) + await expect(isModelSafeWorkspaceFileKey('unrecorded-key')).resolves.toBe(true) }) it('still lets an unknown contribution dominate an unrecorded one', () => { 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 6af92cabbe0..69be5efc158 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 @@ -15,6 +15,8 @@ import { } from '@/lib/execution/durable-secret-provenance' import { reportDurableSecretProvenanceRefusal, + reportDurableSecretProvenanceUnrecorded, + reportDurableSecretProvenanceUnrecordedBatch, reportDurableSecretProvenanceWrite, } from '@/lib/execution/durable-secret-provenance-telemetry' import { @@ -38,8 +40,8 @@ export const MODEL_UNSAFE_WORKSPACE_FILE_ERROR_MESSAGE = /** * Exact sidecars identify known secret contributions. Unrecorded sidecars represent missing - * writer provenance; unknown sidecars represent taint or an invalid binding. Both refuse reads - * for tracked files, while null tracking markers retain legacy compatibility. + * writer provenance and retain ordinary-input behavior with an audit of the coverage gap. Unknown + * sidecars represent taint or an invalid binding and refuse reads. Null markers remain compatible. */ export type WorkspaceFileSecretProvenance = | { status: 'exact'; entries: readonly WorkspaceFileSecretProvenanceEntry[] } @@ -115,9 +117,8 @@ interface ModelSafeWorkspaceFileRow { /** * Combines byte-contributing classifications without broadening any source. * - * An absence stays unrecorded only when no contributor carries known secrets. Mixing it with - * known secret entries cannot preserve an exact classification or discard those entries into a - * permissive absence; the newly combined bytes must remain unknown. + * Missing evidence contributes no known entries. Preserve every known contribution when an + * unrecorded source is combined with an exact source; protection failures still dominate. */ export function mergeWorkspaceFileSecretProvenance( ...provenances: readonly WorkspaceFileSecretProvenance[] @@ -125,13 +126,7 @@ export function mergeWorkspaceFileSecretProvenance( if (provenances.some((provenance) => provenance.status === 'unknown')) { return { status: 'unknown' } } - if (provenances.some((provenance) => provenance.status === 'unrecorded')) { - return provenances.some( - (provenance) => provenance.status === 'exact' && provenance.entries.length > 0 - ) - ? { status: 'unknown' } - : { status: 'unrecorded' } - } + const unrecorded = provenances.some((provenance) => provenance.status === 'unrecorded') const entries = new Map() let bytes = 0 @@ -161,7 +156,9 @@ export function mergeWorkspaceFileSecretProvenance( entries.set(key, entry) } } - return { status: 'exact', entries: [...entries.values()] } + return unrecorded && entries.size === 0 + ? { status: 'unrecorded' } + : { status: 'exact', entries: [...entries.values()] } } function exactEntryByteSize(entry: WorkspaceFileSecretProvenanceEntry): number { @@ -265,7 +262,7 @@ export async function createWorkspaceFileSecretProvenanceFromRegistry( representations: readonly WorkspaceFileSecretProvenanceRepresentation[] = [], representationsComplete = true ): Promise { - /** Missing recording context is persisted separately from known taint; both refuse later reads. */ + /** Missing recording context is persisted separately from known taint. */ if (!registry) return { safe: true, provenance: { status: 'unrecorded' } } const sourceProvenance = registry.exportCommittedProvenanceForValue(sourceValue) const persistedProvenance = Object.is(sourceValue, persistedValue) @@ -387,8 +384,15 @@ export async function createWorkspaceFileSecretProvenanceFromRegistry( } } -function isValidStoredEntries(value: unknown): value is StoredWorkspaceFileSecretProvenanceEntry[] { - if (!Array.isArray(value) || value.length > PROVENANCE_MAX_ENTRIES) { +function isValidStoredEntries( + value: unknown, + status: string | null +): value is StoredWorkspaceFileSecretProvenanceEntry[] { + if ( + !Array.isArray(value) || + value.length > PROVENANCE_MAX_ENTRIES || + ((status === null || status === 'unrecorded') && value.length > 0) + ) { return false } let bytes = 0 @@ -607,7 +611,7 @@ export async function snapshotWorkspaceFileSecretProvenanceInTx( if ( !stored || stored.contentUpdatedAt.getTime() !== contentUpdatedAt.getTime() || - !isValidStoredEntries(stored.entries) + !isValidStoredEntries(stored.entries, stored.status) ) { return { status: 'unknown', entries: [] } } @@ -616,6 +620,16 @@ export async function snapshotWorkspaceFileSecretProvenanceInTx( return { status: 'unknown', entries: [] } } +/** Decodes a captured file revision without dropping anonymous entries or accepting malformed absence. */ +export function workspaceFileSecretProvenanceFromSnapshot( + snapshot: WorkspaceFileSecretProvenanceSnapshot +): WorkspaceFileSecretProvenance { + if (!isValidStoredEntries(snapshot.entries, snapshot.status)) return { status: 'unknown' } + if (snapshot.status === null) return EXACT_EMPTY_WORKSPACE_FILE_SECRET_PROVENANCE + if (snapshot.status !== 'exact') return { status: snapshot.status } + return { status: 'exact', entries: deserializeExactEntriesFromStorage(snapshot.entries) } +} + /** * Binds a previously captured snapshot to the file's new content version, writing the stored entries * verbatim so a revert restores exactly the classification its bytes carried. An untracked snapshot @@ -627,15 +641,12 @@ async function reinstateWorkspaceFileSecretProvenanceInTx( contentUpdatedAt: Date, snapshot: WorkspaceFileSecretProvenanceSnapshot ): Promise { - const provenance: WorkspaceFileSecretProvenance = - snapshot.status === null - ? EXACT_EMPTY_WORKSPACE_FILE_SECRET_PROVENANCE - : snapshot.status !== 'exact' - ? { status: snapshot.status } - : isValidStoredEntries(snapshot.entries) - ? { status: 'exact', entries: deserializeExactEntriesFromStorage(snapshot.entries) } - : { status: 'unknown' } - await replaceWorkspaceFileSecretProvenanceInTx(tx, fileId, contentUpdatedAt, provenance) + await replaceWorkspaceFileSecretProvenanceInTx( + tx, + fileId, + contentUpdatedAt, + workspaceFileSecretProvenanceFromSnapshot(snapshot) + ) } /** @@ -753,7 +764,7 @@ export async function copyWorkspaceFileSecretProvenanceInTx( } if ( source.provenanceContentUpdatedAt?.getTime() !== source.fileContentUpdatedAt.getTime() || - !isValidStoredEntries(source.entries) + !isValidStoredEntries(source.entries, source.status) ) { await replaceWorkspaceFileSecretProvenanceInTx(tx, targetFileId, target.contentUpdatedAt, { status: 'unknown', @@ -895,10 +906,10 @@ export async function getBoundWorkspaceFileSecretProvenance( if (row.secretProvenanceVersion !== 1) return { status: 'unknown' } const bindingIsCurrent = row.provenanceContentUpdatedAt?.getTime() === row.fileContentUpdatedAt.getTime() - if (!bindingIsCurrent || !isValidStoredEntries(row.entries)) return { status: 'unknown' } + if (!bindingIsCurrent || !isValidStoredEntries(row.entries, row.status)) + return { status: 'unknown' } /** - * Preserve the writer's absence/taint distinction for diagnostics. Readers refuse both; - * legacy compatibility is determined only by the tracking marker above. + * Preserve absence separately from protection failures so consumers apply the matching policy. */ if (row.status === 'unrecorded') return { status: 'unrecorded' } if (row.status !== 'exact') return { status: 'unknown' } @@ -950,7 +961,7 @@ export async function getBoundWorkspaceFileSecretProvenanceByMetadata( if ( row.secretProvenanceVersion !== 1 || row.provenanceContentUpdatedAt?.getTime() !== row.fileContentUpdatedAt.getTime() || - !isValidStoredEntries(row.entries) + !isValidStoredEntries(row.entries, row.status) ) { result.set(row.id, { status: 'unknown' }) continue @@ -975,14 +986,29 @@ export async function getBoundWorkspaceFileSecretProvenanceByMetadata( return result } -/** Refuses tracked content whose writer could not record its provenance. */ -function refuseUnrecordedWorkspaceFile(workspaceId: string | undefined): false { - reportDurableSecretProvenanceRefusal({ +/** Missing producer evidence does not prevent ordinary file use; record the coverage gap. */ +function allowUnrecordedWorkspaceFile( + workspaceId: string | undefined, + resourceId?: string, + actorUserId?: string +): true { + reportDurableSecretProvenanceUnrecorded({ surface: 'workspace-file', - cause: 'workspace-file-unrecorded-enforced', workspaceId, + resourceId, + actorUserId, }) - return false + return true +} + +function reportUnrecordedWorkspaceFileBatch(counts: ReadonlyMap): void { + reportDurableSecretProvenanceUnrecordedBatch( + Array.from(counts, ([workspaceId, recordCount]) => ({ + surface: 'workspace-file', + workspaceId, + recordCount, + })) + ) } /** Reports the canonical file identity without exposing its storage key or contents. */ @@ -1006,8 +1032,8 @@ function refuseWorkspaceFileProvenance( /** * Authorizes one model-facing view of an exact workspace-file version. Complete text views import * the entire sidecar so representation-changing consumers retain the original lineage. Derived - * text views import only entries present in the returned value; opaque bytes cannot be inspected - * and therefore require an exact-empty sidecar. + * text views import only entries present in the returned value. Opaque bytes refuse known secret + * contributions; explicit missing producer evidence retains the ordinary-input policy. */ export async function importWorkspaceFileSecretProvenanceForModelView(args: { workspaceId: string @@ -1025,7 +1051,7 @@ export async function importWorkspaceFileSecretProvenanceForModelView(args: { ) } if (provenance.status === 'unrecorded') { - return refuseUnrecordedWorkspaceFile(args.workspaceId) + return allowUnrecordedWorkspaceFile(args.workspaceId, args.identity.fileId) } if (provenance.entries.length === 0) return true if (args.view === 'opaque' || !args.registry) { @@ -1053,7 +1079,7 @@ export async function importWorkspaceFileSecretProvenanceForModelView(args: { ) } -/** Allows opaque bytes to leave private storage only when their exact sidecar is provably empty. */ +/** Admits opaque bytes without known secret contributions or a protection failure. */ export async function isOpaqueWorkspaceFileEgressSafe( workspaceId: string, identity: WorkspaceFileSecretProvenanceIdentity @@ -1066,7 +1092,8 @@ export async function isOpaqueWorkspaceFileEgressSafe( identity.fileId ) } - if (provenance.status === 'unrecorded') return refuseUnrecordedWorkspaceFile(workspaceId) + if (provenance.status === 'unrecorded') + return allowUnrecordedWorkspaceFile(workspaceId, identity.fileId) return ( provenance.entries.length === 0 || refuseWorkspaceFileProvenance( @@ -1112,7 +1139,7 @@ export async function importWorkspaceFileSnapshotProvenance(args: { ) } if (provenance.status === 'unrecorded') { - return refuseUnrecordedWorkspaceFile(args.workspaceId) + return allowUnrecordedWorkspaceFile(args.workspaceId, args.resourceId, args.actorUserId) } if (provenance.entries.length === 0) return true if (!args.registry) { @@ -1160,7 +1187,7 @@ export async function filterModelSafeWorkspaceFileAttachments< const rowByKey = new Map(rows.map((row) => [row.key, row])) const versionKeys = await findWorkspaceFileVersionKeys(keys.filter((key) => !rowByKey.has(key))) - let unrecorded = 0 + const unrecordedByWorkspace = new Map() let refused = 0 const kept = attachments.filter((attachment) => { if (typeof attachment.key !== 'string' || attachment.key.length === 0) return true @@ -1183,20 +1210,18 @@ export async function filterModelSafeWorkspaceFileAttachments< refused += 1 return false } - unrecorded += 1 - return false + const workspaceId = options.workspaceId ?? row.workspaceId ?? undefined + unrecordedByWorkspace.set(workspaceId, (unrecordedByWorkspace.get(workspaceId) ?? 0) + 1) + return true }) if (refused > 0) { refuseWorkspaceFileProvenance('workspace-file-provenance-unavailable', options.workspaceId) } - /** One report for the whole set of attachments, which is one read, rather than one per file. */ - if (unrecorded > 0) { - refuseUnrecordedWorkspaceFile(options.workspaceId) - } + reportUnrecordedWorkspaceFileBatch(unrecordedByWorkspace) return kept } -/** Keeps missing writer provenance distinct from taint for refusal diagnostics. */ +/** Keeps missing writer provenance distinct from a protection failure. */ type ModelSafeWorkspaceFileClassification = 'safe' | 'unrecorded' | 'unsafe' function classifyModelSafeWorkspaceFileRow( @@ -1208,7 +1233,7 @@ function classifyModelSafeWorkspaceFileRow( if (row.secretProvenanceVersion !== 1) return 'unsafe' const bindingIsCurrent = row.provenanceContentUpdatedAt?.getTime() === row.fileContentUpdatedAt.getTime() - if (!bindingIsCurrent || !isValidStoredEntries(row.entries)) return 'unsafe' + if (!bindingIsCurrent || !isValidStoredEntries(row.entries, row.status)) return 'unsafe' if (row.status === 'unrecorded') return 'unrecorded' if (row.status !== 'exact') return 'unsafe' return row.entries.length === 0 ? 'safe' : 'unsafe' @@ -1287,7 +1312,7 @@ export async function areModelSafeWorkspaceFileKeys( ) } - let unrecorded = 0 + const unrecordedByWorkspace = new Map() for (const row of rows) { if ( row.context !== 'workspace' && @@ -1303,8 +1328,11 @@ export async function areModelSafeWorkspaceFileKeys( options.workspaceId ?? row.workspaceId ?? undefined ) } - if (classification === 'unrecorded') unrecorded += 1 + if (classification === 'unrecorded') { + const workspaceId = options.workspaceId ?? row.workspaceId ?? undefined + unrecordedByWorkspace.set(workspaceId, (unrecordedByWorkspace.get(workspaceId) ?? 0) + 1) + } } - /** One report for the batch, not one per key: a caller checking many keys is one read. */ - return unrecorded === 0 || refuseUnrecordedWorkspaceFile(options.workspaceId) + reportUnrecordedWorkspaceFileBatch(unrecordedByWorkspace) + return true } diff --git a/apps/sim/lib/uploads/contexts/workspace/workspace-file-versions.ts b/apps/sim/lib/uploads/contexts/workspace/workspace-file-versions.ts index 9d4a400525e..6e8b0dbc743 100644 --- a/apps/sim/lib/uploads/contexts/workspace/workspace-file-versions.ts +++ b/apps/sim/lib/uploads/contexts/workspace/workspace-file-versions.ts @@ -657,11 +657,13 @@ export function getWorkspaceFileVersion( /** * The provenance captured with one stored version's bytes, which a revert reinstates; null when the - * row is gone, so a caller never reinstates a classification it did not read. + * row is gone or no longer matches captured bytes. Reverts can use the version alone; downloads + * bind the classification to the storage key they opened. */ export async function getWorkspaceFileVersionProvenance( fileId: string, - version: number + version: number, + expectedKey?: string ): Promise { const [row] = await db .select({ @@ -669,7 +671,13 @@ export async function getWorkspaceFileVersionProvenance( entries: workspaceFileVersion.secretProvenanceEntries, }) .from(workspaceFileVersion) - .where(and(eq(workspaceFileVersion.fileId, fileId), eq(workspaceFileVersion.version, version))) + .where( + and( + eq(workspaceFileVersion.fileId, fileId), + eq(workspaceFileVersion.version, version), + expectedKey !== undefined ? eq(workspaceFileVersion.key, expectedKey) : undefined + ) + ) .limit(1) if (!row) return null return { status: toSnapshotStatus(row.status), entries: row.entries } diff --git a/apps/sim/lib/uploads/upload-session/workspace-file-provenance.test.ts b/apps/sim/lib/uploads/upload-session/workspace-file-provenance.test.ts index 2abddd9c130..0213730c9f7 100644 --- a/apps/sim/lib/uploads/upload-session/workspace-file-provenance.test.ts +++ b/apps/sim/lib/uploads/upload-session/workspace-file-provenance.test.ts @@ -46,6 +46,8 @@ describe('private workspace-upload classification', () => { { ...bound, version: 2 }, { ...bound, workspaceId: 'other-workspace' }, { ...bound, provenance: { status: 'future-status' } }, + { ...bound, provenance: { status: 'unrecorded', entries: [entry] } }, + { ...bound, provenance: { status: 'unrecorded', entries: {} } }, { ...bound, provenance: { status: 'exact', entries: {} } }, { ...bound, diff --git a/apps/sim/lib/uploads/upload-session/workspace-file-provenance.ts b/apps/sim/lib/uploads/upload-session/workspace-file-provenance.ts index c7b231dd27b..f47532585a9 100644 --- a/apps/sim/lib/uploads/upload-session/workspace-file-provenance.ts +++ b/apps/sim/lib/uploads/upload-session/workspace-file-provenance.ts @@ -48,7 +48,12 @@ export function readWorkspaceFileUploadProvenance(session: { export function parseWorkspaceFileSecretProvenance(value: unknown): WorkspaceFileSecretProvenance { if (!value || typeof value !== 'object' || !('status' in value)) return { status: 'unknown' } - if (value.status === 'unknown' || value.status === 'unrecorded') return { status: value.status } + if (value.status === 'unknown') return { status: 'unknown' } + if (value.status === 'unrecorded') { + return !('entries' in value) || (Array.isArray(value.entries) && value.entries.length === 0) + ? { status: 'unrecorded' } + : { status: 'unknown' } + } if (value.status !== 'exact' || !('entries' in value)) return { status: 'unknown' } const normalized = normalizeDurableSecretProvenanceEntries(value.entries) if (!normalized) return { status: 'unknown' } diff --git a/apps/sim/lib/workspace-files/application/file-versions.ts b/apps/sim/lib/workspace-files/application/file-versions.ts index edf30250e41..1b7cd7da919 100644 --- a/apps/sim/lib/workspace-files/application/file-versions.ts +++ b/apps/sim/lib/workspace-files/application/file-versions.ts @@ -14,6 +14,11 @@ import { updateWorkspaceFileContent, type WorkspaceFileRecord, } from '@/lib/uploads/contexts/workspace' +import { + getBoundWorkspaceFileSecretProvenance, + type WorkspaceFileSecretProvenance, + workspaceFileSecretProvenanceFromSnapshot, +} from '@/lib/uploads/contexts/workspace/workspace-file-secret-provenance' import { getCurrentWorkspaceFileVersion, getWorkspaceFileVersion, @@ -29,6 +34,10 @@ import { type DownloadWorkspaceFileStreamResult, streamWorkspaceFileRecord, } from '@/lib/workspace-files/application/download-workspace-file' +import { + hasWorkspaceFileDeliveryObserver, + reportWorkspaceFileDelivery, +} from '@/lib/workspace-files/application/file-delivery-observer' import { parseWorkspaceFileRevision } from '@/lib/workspace-files/application/file-revision' import { resolveWorkspaceFileVersionWrite } from '@/lib/workspace-files/application/file-version-write' import { fileOperations } from '@/lib/workspace-files/application/operations' @@ -110,6 +119,23 @@ async function loadVersion( return record } +async function readVersionSecretProvenance( + file: WorkspaceFileRecord, + version: WorkspaceFileVersionRecord +): Promise { + if (!hasWorkspaceFileDeliveryObserver()) return undefined + if (version.isCurrent) { + return getBoundWorkspaceFileSecretProvenance(file.workspaceId, { + fileId: file.id, + key: file.key, + context: file.storageContext ?? 'workspace', + contentUpdatedAt: file.contentUpdatedAt ?? undefined, + }) + } + const snapshot = await getWorkspaceFileVersionProvenance(file.id, version.version, version.key) + return snapshot ? workspaceFileSecretProvenanceFromSnapshot(snapshot) : { status: 'unknown' } +} + /** * Runs a read of a version's stored object, answering 404 when the object is gone — retention or a * delete can remove a superseded version between loading its row and reading its bytes. @@ -182,8 +208,15 @@ export const readWorkspaceFileVersionText = defineAuthorizedWorkspaceFileUseCase const file = await loadActiveFile(context) const version = await loadVersion(file, input.version) const fileAtVersion = recordAtVersion(file, version) + const secretProvenance = await readVersionSecretProvenance(file, version) const result = await readVersionObject(version.version, () => - extractWorkspaceFileRecordText(fileAtVersion, input, principal, request?.signal) + extractWorkspaceFileRecordText( + fileAtVersion, + input, + principal, + request?.signal, + secretProvenance + ) ) return { ...result, file: fileAtVersion, version } }, @@ -196,6 +229,7 @@ export const downloadWorkspaceFileVersion = defineAuthorizedWorkspaceFileUseCase async execute({ input, context, principal }): Promise { const file = await loadActiveFile(context) const version = await loadVersion(file, input.version) + await reportWorkspaceFileDelivery(await readVersionSecretProvenance(file, version)) const result = await readVersionObject(version.version, () => streamWorkspaceFileRecord(recordAtVersion(file, version), principal) ) diff --git a/apps/sim/lib/workspace-files/application/read-workspace-file-text.ts b/apps/sim/lib/workspace-files/application/read-workspace-file-text.ts index 11ac0c2ee7d..07c66477ab5 100644 --- a/apps/sim/lib/workspace-files/application/read-workspace-file-text.ts +++ b/apps/sim/lib/workspace-files/application/read-workspace-file-text.ts @@ -128,8 +128,21 @@ export async function extractWorkspaceFileRecordText( 'maxBytes' | 'offset' | 'limit' | 'allowPlainText' | 'includeSecretProvenance' >, principal: Principal, - signal?: AbortSignal + signal?: AbortSignal, + sourceProvenance?: WorkspaceFileSecretProvenance ): Promise { + const secretProvenance = + sourceProvenance ?? + (input.includeSecretProvenance || hasWorkspaceFileDeliveryObserver() + ? await getBoundWorkspaceFileSecretProvenance(file.workspaceId, { + fileId: file.id, + key: file.key, + context: file.storageContext ?? 'workspace', + contentUpdatedAt: file.contentUpdatedAt ?? undefined, + }) + : undefined) + await reportWorkspaceFileDelivery(secretProvenance) + const extension = input.allowPlainText ? (workspaceFileTextFormat(file) ?? getFileExtension(file.name)) : getFileExtension(file.name) @@ -162,16 +175,6 @@ export async function extractWorkspaceFileRecordText( : await readSourceBuffer(file, maxBytes, signal) const parsed = await parseFileText(content, extension, file.name, signal) const metadata = parsed.metadata ?? {} - const secretProvenance = - input.includeSecretProvenance || hasWorkspaceFileDeliveryObserver() - ? await getBoundWorkspaceFileSecretProvenance(file.workspaceId, { - fileId: file.id, - key: file.key, - context: file.storageContext ?? 'workspace', - contentUpdatedAt: file.contentUpdatedAt ?? undefined, - }) - : undefined - await reportWorkspaceFileDelivery(secretProvenance) const truncated = metadata.truncated === true const { text, lineRange } = sliceFileTextLines( diff --git a/packages/testing/src/mocks/workspace-file-secret-provenance.mock.ts b/packages/testing/src/mocks/workspace-file-secret-provenance.mock.ts index 7feef1d372d..e0b48379187 100644 --- a/packages/testing/src/mocks/workspace-file-secret-provenance.mock.ts +++ b/packages/testing/src/mocks/workspace-file-secret-provenance.mock.ts @@ -33,13 +33,7 @@ function mergeWorkspaceFileSecretProvenance( if (provenances.some((provenance) => provenance.status === 'unknown')) { return { status: 'unknown' } } - if (provenances.some((provenance) => provenance.status === 'unrecorded')) { - return provenances.some( - (provenance) => provenance.status === 'exact' && provenance.entries.length > 0 - ) - ? { status: 'unknown' } - : { status: 'unrecorded' } - } + const unrecorded = provenances.some((provenance) => provenance.status === 'unrecorded') const entries = new Map() let bytes = 0 for (const provenance of provenances) { @@ -68,7 +62,9 @@ function mergeWorkspaceFileSecretProvenance( entries.set(key, entry) } } - return { status: 'exact', entries: [...entries.values()] } + return unrecorded && entries.size === 0 + ? { status: 'unrecorded' } + : { status: 'exact', entries: [...entries.values()] } } /** @@ -94,6 +90,7 @@ export const workspaceFileSecretProvenanceMockFns = { mockInitializeWorkspaceFileSecretProvenanceInTx: vi.fn(), mockPreserveWorkspaceFileSecretProvenanceInTx: vi.fn(), mockSnapshotWorkspaceFileSecretProvenanceInTx: vi.fn(), + mockWorkspaceFileSecretProvenanceFromSnapshot: vi.fn(), mockApplyWorkspaceFileSecretProvenancePolicyInTx: vi.fn(), mockCopyWorkspaceFileSecretProvenanceInTx: vi.fn(), mockMarkWorkspaceFileSecretProvenanceUnknown: vi.fn(), @@ -137,6 +134,7 @@ export const workspaceFileSecretProvenanceMock = { initializeWorkspaceFileSecretProvenanceInTx: fns.mockInitializeWorkspaceFileSecretProvenanceInTx, preserveWorkspaceFileSecretProvenanceInTx: fns.mockPreserveWorkspaceFileSecretProvenanceInTx, snapshotWorkspaceFileSecretProvenanceInTx: fns.mockSnapshotWorkspaceFileSecretProvenanceInTx, + workspaceFileSecretProvenanceFromSnapshot: fns.mockWorkspaceFileSecretProvenanceFromSnapshot, applyWorkspaceFileSecretProvenancePolicyInTx: fns.mockApplyWorkspaceFileSecretProvenancePolicyInTx, copyWorkspaceFileSecretProvenanceInTx: fns.mockCopyWorkspaceFileSecretProvenanceInTx,