diff --git a/.github/workflows/test-build.yml b/.github/workflows/test-build.yml index 54659595287..4c816306b99 100644 --- a/.github/workflows/test-build.yml +++ b/.github/workflows/test-build.yml @@ -129,6 +129,24 @@ jobs: if-no-files-found: warn retention-days: 14 + - name: Verify Google document reads over real HTTP + if: matrix.provision == 'push' + working-directory: apps/sim + env: + NEXT_PUBLIC_APP_URL: http://127.0.0.1:3040 + NEXT_PUBLIC_FORCE_HOSTED: 'false' + SEARCH_GOOGLE_CONTENT_REPORT_PATH: ${{ runner.temp }}/search-google-content.json + run: bun scripts/test-search-google-content-e2e.ts + + - name: Upload Google content acceptance report + if: failure() && matrix.provision == 'push' + uses: actions/upload-artifact@ea165f8d65b6e75b540449e92b4886f43607fa02 # v4 + with: + name: search-google-content + path: ${{ runner.temp }}/search-google-content.json + if-no-files-found: ignore + retention-days: 7 + - name: Verify SCIM and administration over real HTTP working-directory: apps/sim env: diff --git a/apps/sim/lib/file-parsers/docx-parser.ts b/apps/sim/lib/file-parsers/docx-parser.ts index 35586d16b6b..e5c0ba73a15 100644 --- a/apps/sim/lib/file-parsers/docx-parser.ts +++ b/apps/sim/lib/file-parsers/docx-parser.ts @@ -1,5 +1,6 @@ import { readFile } from 'fs/promises' import { createLogger } from '@sim/logger' +import { isRecordLike, toRecord } from '@sim/utils/object' import mammoth from 'mammoth' import { FileParserError, @@ -19,6 +20,79 @@ import { assertOoxmlArchiveWithinLimits } from '@/lib/file-parsers/zip-guard' const logger = createLogger('DocxParser') +/** Bounds repeated notes and generated markup independently of the normalized text budget. */ +const MAX_DOCX_CONVERSION_NODES = 50_000 +const MAX_DOCX_CONVERSION_BYTES = 2 * 1024 * 1024 + +/** + * Mammoth's supported transform hook runs before HTML generation. Count every reference + * expansion, including repeated notes, without retaining the expanded graph. The node ceiling + * also stops cyclic references. The separate 2 MiB ceiling charges raw UTF-8 model strings, + * including attributes, before HTML escaping. Fixed default styles and omitted image data + * bound conversion amplification; embedded style maps could add arbitrary wrapper markup. + */ +function assertDocxConversionWithinLimits(document: unknown, signal?: AbortSignal): void { + const notes = toRecord(toRecord(document).notes) + const pending: unknown[] = [document] + let nodes = 0 + let bytes = 0 + const complexity = () => + new FileParserError('complexity_limit', 'DOCX conversion exceeds its graph budget') + const append = (children: unknown) => { + if (!Array.isArray(children)) + throw new FileParserError('invalid_format', 'DOCX conversion has invalid children') + if (nodes + pending.length + children.length > MAX_DOCX_CONVERSION_NODES) throw complexity() + for (let index = children.length - 1; index >= 0; index--) pending.push(children[index]) + } + while (pending.length) { + signal?.throwIfAborted() + const node = pending.pop() + if (!isRecordLike(node) || typeof node.type !== 'string') + throw new FileParserError('invalid_format', 'DOCX conversion has an invalid node') + if (++nodes > MAX_DOCX_CONVERSION_NODES) throw complexity() + for (const value of Object.values(node)) { + if (typeof value === 'string') bytes += Buffer.byteLength(value, 'utf8') + } + if (bytes > MAX_DOCX_CONVERSION_BYTES) throw complexity() + switch (node.type) { + case 'document': + case 'paragraph': + case 'run': + case 'hyperlink': + case 'table': + case 'tableRow': + case 'tableCell': + append(node.children) + break + case 'note': + case 'comment': + append(node.body) + break + case 'noteReference': { + if (typeof notes.resolve !== 'function') + throw new FileParserError('invalid_format', 'DOCX note resolver is unavailable') + const note: unknown = Reflect.apply(notes.resolve, notes, [node]) + if (!note) throw new FileParserError('invalid_format', 'DOCX references a missing note') + append([note]) + break + } + case 'text': + if (typeof node.value !== 'string') + throw new FileParserError('invalid_format', 'DOCX text is malformed') + break + case 'image': + case 'tab': + case 'checkbox': + case 'break': + case 'bookmarkStart': + case 'commentReference': + break + default: + throw new FileParserError('invalid_format', 'DOCX conversion has an unsupported node') + } + } +} + /** * Extracts DOCX text by rendering the document to HTML with mammoth and walking * that HTML with the shared structured-text walker. mammoth's HTML keeps the @@ -45,18 +119,39 @@ export class DocxParser implements FileParser { } assertOoxmlArchiveWithinLimits(buffer) + const maxTextBytes = + options.docxTextMode === 'complete' + ? (options.maxTextBytes ?? MAX_DOCX_CONVERSION_BYTES) + : undefined + if (maxTextBytes !== undefined && (!Number.isSafeInteger(maxTextBytes) || maxTextBytes <= 0)) + throw new FileParserError('complexity_limit', 'Invalid DOCX text byte budget') const extractionErrors: unknown[] = [] let parserReturnedEmpty = false try { - const htmlResult = await mammoth.convertToHtml({ buffer }) + const htmlResult = await mammoth.convertToHtml( + { buffer }, + maxTextBytes === undefined + ? undefined + : { + includeEmbeddedStyleMap: false, + convertImage: mammoth.images.imgElement(async () => ({ src: '' })), + transformDocument: (document: unknown) => { + assertDocxConversionWithinLimits(document, options.signal) + return document + }, + } + ) options.signal?.throwIfAborted() - const structured = this.structuredTextFromHtml(htmlResult.value) + const structured = this.structuredTextFromHtml(htmlResult.value, maxTextBytes !== undefined) if (structured) { + const content = sanitizeTextForUTF8(structured) + if (maxTextBytes !== undefined && Buffer.byteLength(content, 'utf8') > maxTextBytes) + throw new FileParserError('complexity_limit', 'DOCX text exceeds its byte budget') return { - content: sanitizeTextForUTF8(structured), + content, metadata: { extractionMethod: 'mammoth-html', messages: htmlResult.messages, @@ -64,6 +159,9 @@ export class DocxParser implements FileParser { } } + if (maxTextBytes !== undefined) + throw new FileParserError('no_extractable_text', 'No complete DOCX text was extracted') + const rawResult = await mammoth.extractRawText({ buffer }) options.signal?.throwIfAborted() @@ -79,6 +177,7 @@ export class DocxParser implements FileParser { parserReturnedEmpty = true } catch (mammothError) { options.signal?.throwIfAborted() + if (maxTextBytes !== undefined) throw mammothError logger.warn('mammoth failed, trying officeparser:', mammothError) extractionErrors.push(mammothError) } @@ -161,14 +260,16 @@ export class DocxParser implements FileParser { /** * Walks mammoth's HTML rendering under the HTML parser's size caps. A rendering * too large to walk safely falls back to the raw-text path by returning empty, - * since mammoth has already materialised the document once at that point. + * since mammoth has already materialised the document once at that point. Budgeted reads + * reject instead: the fallback cannot preserve the conversion and completeness bounds. */ - private structuredTextFromHtml(html: string): string { + private structuredTextFromHtml(html: string, bounded = false): string { if (!html || html.trim().length === 0) return '' try { assertHtmlStringWithinLimits(html) } catch (error) { if (isHtmlComplexityError(error)) { + if (bounded) throw error logger.warn('mammoth HTML exceeds walker limits, using raw text:', error.message) return '' } diff --git a/apps/sim/lib/file-parsers/pdf-parser.ts b/apps/sim/lib/file-parsers/pdf-parser.ts index 3ae6f213783..a599f3d7806 100644 --- a/apps/sim/lib/file-parsers/pdf-parser.ts +++ b/apps/sim/lib/file-parsers/pdf-parser.ts @@ -36,6 +36,9 @@ export const MAX_PDF_TEXT_CHARS = 10_000_000 /** Complete extraction shares the ingestion pipeline's bounded text-output envelope. */ export const MAX_COMPLETE_PDF_TEXT_BYTES = 20 * 1024 * 1024 +/** Independent retained-text ceiling, including conservative layout overhead. */ +const MAX_COMPLETE_PDF_RETAINED_TEXT_BYTES = MAX_COMPLETE_PDF_TEXT_BYTES + /** Bounds expansion on one page independently of a long document's output budget. */ export const MAX_COMPLETE_PDF_PAGE_CHARS = 250_000 @@ -43,10 +46,8 @@ export const MAX_COMPLETE_PDF_PAGE_CHARS = 250_000 const PDF_EXTRACTION_TIMEOUT_MS = 60_000 /** - * Upper bound on what line reconstruction adds per line after the budget is - * spent: a two-character paragraph break plus a three-character heading marker. - * The complete-mode byte ceiling therefore trips slightly earlier than it did - * when pages were flattened to one line, by at most this many bytes per line. + * Conservative retained-text allowance for a paragraph break and heading marker. + * This protects parser state independently of the caller's rendered-output budget. */ const MAX_LINE_DECORATION_BYTES = 5 @@ -249,8 +250,8 @@ function readPageHeight(page: PDFPageProxy): number | undefined { } } -/** Bytes a page's lines can occupy in the output once joined and decorated. */ -function estimatePageBytes(lines: readonly PdfLine[]): number { +/** Conservative text storage for retained lines and their possible decorations. */ +function estimateRetainedPageBytes(lines: readonly PdfLine[]): number { let bytes = 0 for (const line of lines) { bytes += Buffer.byteLength(line.text, 'utf8') + MAX_LINE_DECORATION_BYTES @@ -267,7 +268,8 @@ function estimatePageBytes(lines: readonly PdfLine[]): number { async function assemblePages( pages: readonly PdfPageLines[], complete: boolean, - signal: AbortSignal | undefined + signal: AbortSignal | undefined, + maxTextBytes: number ): Promise { signal?.throwIfAborted() const filteredPages = suppressFurniture(pages) @@ -280,6 +282,7 @@ async function assemblePages( headingMarkers: PDF_HEADING_MARKERS_ENABLED && headingMarkersViable(allLines, bodyHeight), } const pageTexts: string[] = [] + let outputBytes = 0 for (const [index, lines] of filteredPages.entries()) { if (index > 0 && index % ASSEMBLY_YIELD_EVERY_PAGES === 0) { await sleep(0) @@ -287,7 +290,17 @@ async function assemblePages( } const joined = joinLines(lines, options) const text = complete ? normalizePdfWhitespace(sanitizeTextForUTF8(joined)).trim() : joined - if (text.length > 0) pageTexts.push(text) + if (text.length > 0) { + if (complete) { + outputBytes += + Buffer.byteLength(text, 'utf8') + (pageTexts.length > 0 ? PAGE_SEPARATOR.length : 0) + if (outputBytes > maxTextBytes) + throw completeExtractionLimit( + `PDF text exceeds the safe ${maxTextBytes.toLocaleString()}-byte output limit.` + ) + } + pageTexts.push(text) + } } const text = pageTexts.join(PAGE_SEPARATOR) return complete ? text : normalizePdfWhitespace(text).trim() @@ -297,6 +310,13 @@ function completeExtractionLimit(message: string): FileParserError { return new FileParserError('complexity_limit', `${message} Split or simplify the PDF and retry.`) } +function completeBudget(value: number | undefined, ceiling: number): number { + if (value === undefined) return ceiling + if (!Number.isSafeInteger(value) || value < 1) + throw completeExtractionLimit('PDF extraction limits must be positive safe integers.') + return Math.min(value, ceiling) +} + async function extractTextWithinBudget( pdf: PDFDocumentProxy, options: FileParseOptions, @@ -305,17 +325,21 @@ async function extractTextWithinBudget( const { signal } = options const complete = options.pdfTextMode === 'complete' const totalPages = pdf.numPages - const pageLimit = Math.min(totalPages, MAX_PDF_PAGES) + const maxPages = complete ? completeBudget(options.pdfMaxPages, MAX_PDF_PAGES) : MAX_PDF_PAGES + const maxTextBytes = complete + ? completeBudget(options.maxTextBytes, MAX_COMPLETE_PDF_TEXT_BYTES) + : MAX_COMPLETE_PDF_TEXT_BYTES + const pageLimit = Math.min(totalPages, maxPages) const pages: PdfPageLines[] = [] let remainingChars = MAX_PDF_TEXT_CHARS - let outputBytes = 0 + let retainedTextBytes = 0 let pagesRead = 0 let truncated = totalPages > pageLimit if (complete && truncated) { throw completeExtractionLimit( - `PDF exceeds the safe limit of ${MAX_PDF_PAGES.toLocaleString()} pages.` + `PDF exceeds the safe limit of ${maxPages.toLocaleString()} pages.` ) } @@ -341,13 +365,9 @@ async function extractTextWithinBudget( const page = pageResult const pageHeight = readPageHeight(page) let extraction: PageExtraction + const pageCharLimit = complete ? MAX_COMPLETE_PDF_PAGE_CHARS : remainingChars try { - extraction = await readPageWithinBudget( - page, - complete ? MAX_COMPLETE_PDF_PAGE_CHARS : remainingChars, - deadline, - signal - ) + extraction = await readPageWithinBudget(page, pageCharLimit, deadline, signal) } finally { page.cleanup() } @@ -361,10 +381,11 @@ async function extractTextWithinBudget( pagesRead++ if (complete) { if (lines.length > 0) { - outputBytes += estimatePageBytes(lines) + (pages.length > 0 ? PAGE_SEPARATOR.length : 0) - if (outputBytes > MAX_COMPLETE_PDF_TEXT_BYTES) { + retainedTextBytes += + estimateRetainedPageBytes(lines) + (pages.length > 0 ? PAGE_SEPARATOR.length : 0) + if (retainedTextBytes > MAX_COMPLETE_PDF_RETAINED_TEXT_BYTES) { throw completeExtractionLimit( - `PDF text exceeds the safe ${MAX_COMPLETE_PDF_TEXT_BYTES.toLocaleString()}-byte output limit.` + `PDF text exceeds the safe ${MAX_COMPLETE_PDF_RETAINED_TEXT_BYTES.toLocaleString()}-byte output limit.` ) } pages.push({ lines, pageHeight }) @@ -387,7 +408,7 @@ async function extractTextWithinBudget( } } - let text = await assemblePages(pages, complete, signal) + let text = await assemblePages(pages, complete, signal, maxTextBytes) /** Paragraph breaks land after the budget is spent; trimming that overflow is a truncation too. */ if (!complete && text.length > MAX_PDF_TEXT_CHARS) { diff --git a/apps/sim/lib/file-parsers/types.ts b/apps/sim/lib/file-parsers/types.ts index a55bbab466f..5eb3ab5da33 100644 --- a/apps/sim/lib/file-parsers/types.ts +++ b/apps/sim/lib/file-parsers/types.ts @@ -36,15 +36,22 @@ export interface FileParseOptions { /** Indexing callers require complete extraction; preview row limits must not discard content. */ contentMode?: 'preview' | 'complete' /** - * CSV/XLSX byte budget when contentMode is 'complete' (default 25 MiB). + * Complete text byte budget for CSV/XLSX (default 25 MiB) and PDF (default 20 MiB). + * PDF applies this when pdfTextMode is 'complete' and never exceeds its safe default. * Exceeding it throws complexity_limit instead of returning a truncated prefix. + * DOCX applies this only when docxTextMode is 'complete' (default 2 MiB); separate + * conversion graph limits can reject structurally complex inputs before normalization. * Preview mode retains its own limits; other formats use their parser-specific budgets. */ maxTextBytes?: number /** Preserve textual markup in a canonical .txt artifact instead of interpreting it as HTML or RTF. */ textMode?: 'literal' + /** Opt into bounded DOCX conversion without lossy or unchecked parser fallbacks. */ + docxTextMode?: 'complete' /** Complete PDF extraction rejects safety limits instead of returning preview text. */ pdfTextMode?: 'preview' | 'complete' + /** Lower page ceiling for complete PDF extraction; defaults to the parser's safe limit. */ + pdfMaxPages?: number } export interface FileParser { diff --git a/apps/sim/lib/integrations/availability.server.oauth-configured.test.ts b/apps/sim/lib/integrations/availability.server.oauth-configured.test.ts index c556276a03f..89498c2d0ad 100644 --- a/apps/sim/lib/integrations/availability.server.oauth-configured.test.ts +++ b/apps/sim/lib/integrations/availability.server.oauth-configured.test.ts @@ -1,5 +1,5 @@ -import { setEnv } from '@sim/testing/mocks/env.mock' -import { describe, expect, it } from 'vitest' +import { resetEnvMock, setEnv } from '@sim/testing/mocks/env.mock' +import { afterAll, describe, expect, it } from 'vitest' import { getIntegrationAvailability, getOAuthServiceAvailability, @@ -8,7 +8,10 @@ import { setEnv({ GITHUB_APP_CLIENT_ID: 'repository-client', GITHUB_APP_CLIENT_SECRET: 'repository-secret', + GOOGLE_CLIENT_ID: undefined, + GOOGLE_CLIENT_SECRET: undefined, }) +afterAll(resetEnvMock) describe('OAuth service availability projection', () => { it('uses the repository App while the workflow block keeps its API-key path', () => { diff --git a/apps/sim/lib/integrations/availability.server.test.ts b/apps/sim/lib/integrations/availability.server.test.ts index b3f24e6ade9..64f700324e2 100644 --- a/apps/sim/lib/integrations/availability.server.test.ts +++ b/apps/sim/lib/integrations/availability.server.test.ts @@ -1,5 +1,6 @@ import integrationsJson from '@sim/deployment-config/integrations.json' -import { describe, expect, it } from 'vitest' +import { resetEnvMock, setEnv } from '@sim/testing/mocks/env.mock' +import { afterAll, describe, expect, it } from 'vitest' import { OAUTH_CLIENT_CAPABILITIES, resolveOAuthClientCapabilityId, @@ -23,6 +24,14 @@ import { import type { Integration } from '@/lib/integrations/types' import { getServiceConfigByServiceId } from '@/lib/oauth/utils' +setEnv({ + X_CLIENT_ID: undefined, + X_CLIENT_SECRET: undefined, + GITHUB_APP_CLIENT_ID: undefined, + GITHUB_APP_CLIENT_SECRET: undefined, +}) +afterAll(resetEnvMock) + const integrations = integrationsJson.integrations as readonly Integration[] function availabilityFor( diff --git a/apps/sim/lib/sim-search/live/README.md b/apps/sim/lib/sim-search/live/README.md index 7d39db00629..6dee2df8372 100644 --- a/apps/sim/lib/sim-search/live/README.md +++ b/apps/sim/lib/sim-search/live/README.md @@ -32,6 +32,8 @@ For an explicit source user list, source-side verification tries those users' de The source credential must be able to read the file. Selected folders, accessible subfolders, shared-drive IDs, and file types are checked using source-side metadata. Folder queries narrow candidate retrieval; metadata verification remains authoritative. A member's personal file outside the configured source scope is excluded even if their OAuth token can read it. Folder expansion and ancestry traversal are bounded and may report partial coverage. +Text-bearing PDF and DOCX files are downloaded with the member credential after checking `capabilities.canDownload`, then parsed with the shared document parsers. The download is capped at 4 MiB of decoded bytes; PDF extraction is complete or unavailable, with a 100-page and 1 MiB text limit. DOCX archives are checked before parsing (10 MiB total expanded, 4 MiB per entry), with a 1 MiB extracted-text limit. Complete DOCX reads additionally bound the expanded conversion graph to 50,000 nodes and 2 MiB of model strings before HTML rendering, use built-in styles without embedding image bytes, and reject conversion failures instead of retrying a lossy fallback. MIME/signature mismatches, password protection, malformed files, and documents with no extractable text fail explicitly; this path does not perform OCR. Existing source links and bounded comments/replies remain in successful reads. Other unsupported binary formats return labeled metadata only. + ### Gmail Member mode searches the connected mailbox using Gmail message search operators. Service mode requires Workspace delegation and verifies the connected mailbox's email through Gmail, then validates its active Directory identity/customer and any selected source user list. Delegation targets that same mailbox; it never substitutes another user's mailbox. The admin label picker browses the delegated administrator's mailbox but stores label names, since custom label IDs differ between mailboxes. @@ -42,7 +44,9 @@ For a source configured with **Labels = INBOX**: 2. Fetch each candidate's metadata under the source's delegated token for that member's mailbox. 3. Check the message's current `labelIds` against INBOX. Custom labels are resolved in that mailbox; labels may differ across users. 4. Apply the source's rolling date range, Promotions/Social exclusions, and custom Gmail query. A custom query is verified through a source-side message search restricted by the message's RFC Message-ID, followed by exact Gmail message-ID matching. -5. Return only verified messages. Reading a result fetches that individual message through the member's connection and rechecks the source boundary; it does not expose other messages in the thread. +5. Return only verified messages. Reading a result keeps that message as the anchor, derives its thread from a fresh provider response, and fetches thread metadata before selecting at most seven sibling candidates. Every sibling passes the same source boundary before its body is fetched. The anchor and all included siblings are checked again against current source settings before any conversation content is returned; stale labels from the read cannot authorize the final response. + +Conversation output is chronological and attributes each message to its author, timestamp and message URL. Search filters identify the anchor; context may have different timestamps or wording, while organization source restrictions always apply to each message. Each formatted message is capped at 24,000 characters. MIME traversal is capped at 32 levels and 256 parts, chooses one alternative body, and excludes named attachments. Missing external bodies are labeled as previews. The reader marks omitted or truncated context explicitly and does not download attachments. Thread metadata and each body use the shared response-size limit and request budget. Selected labels are alternatives. Source date/query settings are authoritative result checks; pagination may be necessary to find more allowed candidates. Custom-query pages are reduced to fit the additional permission-check requests while preserving Gmail's continuation token. The existing source defaults exclude Promotions and Social unless disabled. Google documents that API search matches messages and differs from Gmail UI thread matching and alias expansion: [Gmail filtering guide](https://developers.google.com/workspace/gmail/api/guides/filtering). @@ -119,7 +123,7 @@ Self-managed GitLab is resolved from the saved source's validated host/project i | Provider | Search | Read | Service boundary | | --- | --- | --- | --- | | Drive | `GET /drive/v3/files` with Drive `q` | File metadata, Docs/Slides exports, Sheets values, or supported text media | Delegated Drive file visibility and configured folders/types | -| Gmail | `GET /gmail/v1/users/me/messages` with Gmail operators | That exact message with `format=full` | Same-mailbox delegation, labels, date range, category and custom-query checks | +| Gmail | `GET /gmail/v1/users/me/messages` with Gmail operators | Anchor plus up to seven independently authorized conversation messages | Same-mailbox delegation, labels, date range, category and custom-query checks | | Calendar | CalendarList then `/calendars/{id}/events` | `/calendars/{id}/events/{eventId}` | Same-user delegation, selected calendars, event window/query | | Slack | `POST /api/assistant.search.context` | `conversations.replies` or `files.info` preview | Member only; Slack enforces the connected user's grant | | Jira | `POST /ex/jira/{cloudId}/rest/api/3/search/jql` | `/rest/api/3/issue/{key}` under that cloud site | Member only | diff --git a/apps/sim/lib/sim-search/live/account-session.ts b/apps/sim/lib/sim-search/live/account-session.ts index b14cb879600..52cfd8e88f5 100644 --- a/apps/sim/lib/sim-search/live/account-session.ts +++ b/apps/sim/lib/sim-search/live/account-session.ts @@ -172,7 +172,13 @@ export async function openLiveAccountSession( }, read(reference, filters) { if (admin) return admin.read(reference) - if (client) return readNativeProvider(provider, client, reference, boundary.policy, filters) + if (client) + return readNativeProvider(provider, client, reference, { + policy: boundary.policy, + filters, + signal, + verify: boundary.verify, + }) return readMcp(reference.id) }, async verifyCurrent(document) { @@ -180,7 +186,13 @@ export async function openLiveAccountSession( livePolicyFor(await loadLiveSearchPolicies(owner), provider), true ) - return current.verify(document) + const { id, container, kind } = document + if (!(await current.verify({ id, container, kind }))) return false + for (const dependency of document.accessDependencies ?? []) { + signal.throwIfAborted() + if (!(await current.verify(dependency))) return false + } + return true }, } } diff --git a/apps/sim/lib/sim-search/live/application.test.ts b/apps/sim/lib/sim-search/live/application.test.ts index 0108c2fa6ac..de633e48bb1 100644 --- a/apps/sim/lib/sim-search/live/application.test.ts +++ b/apps/sim/lib/sim-search/live/application.test.ts @@ -30,7 +30,7 @@ vi.mock('@/lib/sim-search/live/service-session', () => ({ })) vi.mock('@/lib/sim-search/live/http', async (original) => ({ ...(await original()), - createNativeClient: () => ({ json: mocks.json, text: vi.fn() }), + createNativeClient: () => ({ json: mocks.json, text: vi.fn(), bytes: vi.fn() }), })) vi.mock('@/lib/sim-search/live/gitlab-admin', () => ({ createAdminGitLabSession: mocks.admin })) vi.mock('@/lib/sim-search/live/policy-store', () => ({ @@ -71,6 +71,7 @@ import { import { NativeSearchError } from '@/lib/sim-search/live/http' import { createPolicyVerifier } from '@/lib/sim-search/live/policy' import { defaultLiveSearchPolicy } from '@/lib/sim-search/live/policy-schema' +import { livePolicyFor } from '@/lib/sim-search/live/policy-store' import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry' workspaceAuthzMockFns.mockPermissionSatisfies.mockImplementation( @@ -102,6 +103,8 @@ const document = { describe('authorized live retrieval', () => { beforeEach(() => { resetDbChainMock() + vi.mocked(livePolicyFor).mockReset() + vi.mocked(livePolicyFor).mockReturnValue(defaultLiveSearchPolicy()) mocks.service.mockResolvedValue(undefined) knowledgeContextsMockFns.mockResolveKnowledgeOwnerContext.mockResolvedValue({ workspaceId: 'workspace', @@ -365,6 +368,73 @@ describe('authorized live retrieval', () => { ).rejects.toThrow('outside your organization') expect(mocks.read).toHaveBeenCalledOnce() }) + it('suppresses conversation content when a sibling leaves the current service source', async () => { + const gmail = { ...account, provider: 'gmail', providerId: 'gmail', displayName: 'Mail' } + mocks.accounts.mockResolvedValue([gmail]) + mocks.resolveAccount.mockResolvedValue({ account: gmail, accessToken: 'secret' }) + const search = await searchLiveKnowledge.execute({ principal, input }) + mocks.read.mockResolvedValue({ + ...document, + content: 'Anchor evidence. Restricted sibling evidence.', + accessDependencies: [{ id: 'sibling' }], + }) + mocks.service.mockResolvedValueOnce({ + policy: defaultLiveSearchPolicy('gmail'), + verify: async () => true, + partial: false, + }) + mocks.service.mockResolvedValueOnce({ + policy: defaultLiveSearchPolicy('gmail'), + verify: async ({ id }: { id: string }) => id !== 'sibling', + partial: false, + }) + + await expect( + readLiveDocument.execute({ + principal, + input: { + workspaceId: 'workspace', + documentId: search.results[0]!.documentId, + limit: 1, + resultSecretRegistry: new ResolvedSecretTraceRegistry([]), + }, + }) + ).rejects.toThrow('outside your organization') + }) + it('ignores stale message labels when rechecking member access after a read', async () => { + const gmail = { ...account, provider: 'gmail', providerId: 'gmail', displayName: 'Mail' } + mocks.accounts.mockResolvedValue([gmail]) + mocks.resolveAccount.mockResolvedValue({ account: gmail, accessToken: 'secret' }) + const search = await searchLiveKnowledge.execute({ principal, input }) + vi.mocked(livePolicyFor).mockReturnValue({ + ...defaultLiveSearchPolicy('gmail'), + mode: 'selected', + included: ['INBOX'], + }) + let readCompleted = false + mocks.json.mockImplementation(async (path: string) => { + if (path === '/gmail/v1/users/me/labels') return { labels: [{ id: 'INBOX', name: 'INBOX' }] } + if (path === '/gmail/v1/users/me/messages/doc') + return { id: 'doc', labelIds: readCompleted ? ['SENT'] : ['INBOX'] } + throw new Error(`Unexpected Gmail resource: ${path}`) + }) + mocks.read.mockImplementation(async () => { + readCompleted = true + return { ...document, accessMetadata: { id: 'doc', labelIds: ['INBOX'] } } + }) + + await expect( + readLiveDocument.execute({ + principal, + input: { + workspaceId: 'workspace', + documentId: search.results[0]!.documentId, + limit: 1, + resultSecretRegistry: new ResolvedSecretTraceRegistry([]), + }, + }) + ).rejects.toThrow('outside your organization') + }) it('rejects nonmembers before account discovery or provider calls', async () => { workspaceAuthzMockFns.mockResolveEffectiveWorkspacePermission.mockResolvedValue(null) await expect(searchLiveKnowledge.execute({ principal, input })).rejects.toThrow( @@ -434,7 +504,7 @@ describe('authorized live retrieval', () => { verify: createPolicyVerifier( 'google_drive', policy, - { json: mocks.json, text: vi.fn() }, + { json: mocks.json, text: vi.fn(), bytes: vi.fn() }, 'https://www.googleapis.com' ), partial: false, diff --git a/apps/sim/lib/sim-search/live/atlassian.test.ts b/apps/sim/lib/sim-search/live/atlassian.test.ts index fda5330dcdb..69ae3483dee 100644 --- a/apps/sim/lib/sim-search/live/atlassian.test.ts +++ b/apps/sim/lib/sim-search/live/atlassian.test.ts @@ -14,6 +14,7 @@ function client(rows: Record): NativeClient & { json: ReturnTyp return rows[path] }), text: vi.fn(), + bytes: vi.fn(), } } @@ -62,6 +63,7 @@ describe('Confluence live documents', () => { throw new NativeSearchError('unavailable', 'Provider request failed (404).', undefined, 404) }), text: vi.fn(), + bytes: vi.fn(), } await expect(readAtlassian(api, 'confluence', '9', 'cloud')).resolves.toMatchObject({ kind: 'blogpost', @@ -83,6 +85,7 @@ describe('Confluence live documents', () => { throw new Error(`Unexpected request: ${path}`) }), text: vi.fn(), + bytes: vi.fn(), } await expect(readAtlassian(api, 'confluence', '9', 'cloud')).rejects.toBe(failure) }) diff --git a/apps/sim/lib/sim-search/live/coda-service.test.ts b/apps/sim/lib/sim-search/live/coda-service.test.ts index 392c180e4ab..e66306f0a01 100644 --- a/apps/sim/lib/sim-search/live/coda-service.test.ts +++ b/apps/sim/lib/sim-search/live/coda-service.test.ts @@ -5,7 +5,7 @@ import type { NativeClient } from '@/lib/sim-search/live/types' const doc = { id: 'allowed', name: 'Shared doc', browserLink: 'https://coda.io/d/doc' } const reference = { id: 'coda://docs/allowed/pages/page', kind: 'mcp' } -const createClient = () => ({ json: vi.fn(), text: vi.fn() }) +const createClient = () => ({ json: vi.fn(), text: vi.fn(), bytes: vi.fn() }) describe('Coda service document boundary', () => { it('intersects explicit service document selections before loading source metadata', async () => { diff --git a/apps/sim/lib/sim-search/live/coda.test.ts b/apps/sim/lib/sim-search/live/coda.test.ts index a8d4e2a6e2e..b030c041873 100644 --- a/apps/sim/lib/sim-search/live/coda.test.ts +++ b/apps/sim/lib/sim-search/live/coda.test.ts @@ -5,7 +5,7 @@ describe('Coda REST discovery', () => { it('sends only the provider page token on continuation because it encodes the original query', async () => { const json = vi.fn().mockResolvedValue({ items: [] }) await searchCoda( - { json, text: vi.fn() }, + { json, text: vi.fn(), bytes: vi.fn() }, { query: 'launch', limit: 10, diff --git a/apps/sim/lib/sim-search/live/dates.test.ts b/apps/sim/lib/sim-search/live/dates.test.ts index 15636fabd6f..8bfe1337787 100644 --- a/apps/sim/lib/sim-search/live/dates.test.ts +++ b/apps/sim/lib/sim-search/live/dates.test.ts @@ -12,7 +12,7 @@ const filters = { sortBy: 'oldest' as const, } const input: NativeSearchInput = { query: '', limit: 20, scopes: ['search:read.public'], filters } -const client = () => ({ json: vi.fn(), text: vi.fn() }) +const client = () => ({ json: vi.fn(), text: vi.fn(), bytes: vi.fn() }) const doc: NativeDocument = { id: 'one', title: 'Meeting', diff --git a/apps/sim/lib/sim-search/live/discussion-reads.test.ts b/apps/sim/lib/sim-search/live/discussion-reads.test.ts index e588c65f3dc..647f01e6372 100644 --- a/apps/sim/lib/sim-search/live/discussion-reads.test.ts +++ b/apps/sim/lib/sim-search/live/discussion-reads.test.ts @@ -20,12 +20,16 @@ const DRIVE_FILE = { webViewLink: 'https://docs.google.com/document/d/doc/edit', } const text: NativeClient['text'] = async () => 'Original document text' +const bytes: NativeClient['bytes'] = async () => { + throw new Error('Unexpected binary request') +} /** Independent provider wire fixtures exercise pagination, partial failures, and rendering. */ describe('GitHub conversation reads', () => { it('keeps review decisions, diff context and discussion replies separate from the PR body', async () => { const api: NativeClient = { text, + bytes, async json(path) { if (path.endsWith('/issues/42')) return ISSUE if (path.endsWith('/issues/42/comments')) @@ -103,6 +107,7 @@ describe('GitHub conversation reads', () => { it('exposes the matching comment passage in search even when the issue body is nonempty', async () => { const api: NativeClient = { text, + bytes, async json() { return { total_count: 1, @@ -137,6 +142,7 @@ describe('GitHub conversation reads', () => { it('reads subsequent comment pages and keeps the body when a different review endpoint is denied', async () => { const api: NativeClient = { text, + bytes, async json(path, options) { if (path.endsWith('/issues/42')) return ISSUE if (path.endsWith('/issues/42/comments')) @@ -158,6 +164,7 @@ describe('GitHub conversation reads', () => { let requests = 0 const api: NativeClient = { text, + bytes, async json(path, options) { if (++requests > 11) throw new Error('Unbounded discussion pagination') if (path.endsWith('/issues/42')) return ISSUE @@ -178,6 +185,7 @@ describe('GitHub conversation reads', () => { it('marks oversized discussion text incomplete instead of returning an unbounded transcript', async () => { const api: NativeClient = { text, + bytes, async json(path) { if (path.endsWith('/issues/42')) return { ...ISSUE, pull_request: undefined } return [{ id: 1, body: 'x'.repeat(200_000) }] @@ -194,6 +202,7 @@ describe('Drive discussion reads', () => { it('reads all comment pages with full nested replies, author/time context and resolution state', async () => { const api: NativeClient = { text, + bytes, async json(path, options) { if (!path.endsWith('/comments')) return DRIVE_FILE if (options?.query?.pageToken === 'next') @@ -246,6 +255,7 @@ describe('Drive discussion reads', () => { it('reports comments unavailable while preserving readable document content', async () => { const api: NativeClient = { text, + bytes, async json(path) { if (path.endsWith('/comments')) throw new NativeSearchError('rate_limited', 'Provider rate limit reached') @@ -262,6 +272,7 @@ describe('Drive discussion reads', () => { let pages = 0 const api: NativeClient = { text, + bytes, async json(path) { if (!path.endsWith('/comments')) return DRIVE_FILE if (++pages > 3) throw new Error('Unbounded Drive pagination') @@ -281,6 +292,7 @@ describe('Drive discussion reads', () => { it('bounds nested reply content and explicitly marks what was omitted', async () => { const api: NativeClient = { text, + bytes, async json(path) { if (!path.endsWith('/comments')) return DRIVE_FILE return { diff --git a/apps/sim/lib/sim-search/live/drive-content.ts b/apps/sim/lib/sim-search/live/drive-content.ts new file mode 100644 index 00000000000..9c284b97512 --- /dev/null +++ b/apps/sim/lib/sim-search/live/drive-content.ts @@ -0,0 +1,150 @@ +import { toStringOrNull } from '@sim/utils/coerce' +import { toRecord } from '@sim/utils/object' +import { isPayloadSizeLimitError } from '@/lib/core/utils/stream-limits' +import { FileParserError, getFileParserErrorCode } from '@/lib/file-parsers/errors' +import { sniffFileKind } from '@/lib/file-parsers/sniff' +import type { FileParseResult } from '@/lib/file-parsers/types' +import { + assertOoxmlArchiveWithinLimits, + DEFAULT_OOXML_SIZE_LIMITS, +} from '@/lib/file-parsers/zip-guard' +import { NATIVE_RESPONSE_MAX_BYTES, NativeSearchError, segment } from '@/lib/sim-search/live/http' +import type { NativeClient } from '@/lib/sim-search/live/types' + +const MAX_DRIVE_TEXT_BYTES = 1024 * 1024 +const MAX_DRIVE_PDF_PAGES = 100 +const DRIVE_DOCX_LIMITS = { + ...DEFAULT_OOXML_SIZE_LIMITS, + maxTotalUncompressedBytes: 10 * 1024 * 1024, + maxEntryUncompressedBytes: 4 * 1024 * 1024, +} + +/** Extracts bounded text from the exact downloadable PDF or DOCX identified by Drive metadata. */ +export async function readDriveFileContent( + client: NativeClient, + id: string, + row: Record, + signal?: AbortSignal +): Promise { + signal?.throwIfAborted() + if (row.id !== id) + throw new NativeSearchError( + 'unavailable', + 'Drive returned a different file. Search again before reading it.' + ) + if (toRecord(row.capabilities).canDownload !== true) + throw new NativeSearchError( + 'unavailable', + 'Drive does not allow downloading this file. Ask the owner to allow downloads or open the source.' + ) + const mimeType = toStringOrNull(row.mimeType) + const extension = + mimeType === 'application/pdf' + ? 'pdf' + : mimeType === 'application/vnd.openxmlformats-officedocument.wordprocessingml.document' + ? 'docx' + : undefined + if (!extension) + throw new NativeSearchError( + 'unavailable', + 'This file type does not support binary text reads. Open the source.' + ) + if (row.size !== undefined) { + const size = toStringOrNull(row.size) + if (!size || !/^\d{1,20}$/.test(size)) + throw new NativeSearchError( + 'unavailable', + 'Drive returned an invalid file size. Refresh the source and retry.' + ) + if (Number(size) > NATIVE_RESPONSE_MAX_BYTES) + throw new NativeSearchError( + 'unavailable', + 'This file exceeds the 4 MiB live-read limit. Split it into smaller files or open the source.' + ) + } + + let bytes: Buffer + try { + bytes = await client.bytes(`/drive/v3/files/${segment(id)}`, { + alt: 'media', + supportsAllDrives: 'true', + }) + } catch (error) { + signal?.throwIfAborted() + if (isPayloadSizeLimitError(error)) + throw new NativeSearchError( + 'unavailable', + 'This file exceeds the 4 MiB live-read limit. Split it into smaller files or open the source.' + ) + throw error + } + signal?.throwIfAborted() + try { + if (bytes.byteLength === 0) throw new FileParserError('empty_input', 'Empty file') + if (bytes.byteLength > NATIVE_RESPONSE_MAX_BYTES) + throw new FileParserError('complexity_limit', 'File exceeds the binary read limit') + const kind = sniffFileKind(bytes, extension) + if (kind === 'encrypted-ooxml') + throw new FileParserError('encrypted_file', 'Encrypted document') + if (kind !== extension) + throw new FileParserError('invalid_format', 'File content does not match its Drive MIME type') + + let parsed: FileParseResult + if (extension === 'pdf') { + const { PdfParser } = await import('@/lib/file-parsers/pdf-parser') + parsed = await new PdfParser().parseBuffer(bytes, { + signal, + pdfTextMode: 'complete', + pdfMaxPages: MAX_DRIVE_PDF_PAGES, + maxTextBytes: MAX_DRIVE_TEXT_BYTES, + }) + } else { + assertOoxmlArchiveWithinLimits(bytes, DRIVE_DOCX_LIMITS) + const { DocxParser } = await import('@/lib/file-parsers/docx-parser') + parsed = await new DocxParser().parseBuffer(bytes, { + signal, + docxTextMode: 'complete', + maxTextBytes: MAX_DRIVE_TEXT_BYTES, + }) + } + signal?.throwIfAborted() + if (parsed.metadata?.degraded || parsed.metadata?.truncated) + throw new FileParserError('complexity_limit', 'Complete document text was not extracted') + if (!parsed.content.trim()) + throw new FileParserError('no_extractable_text', 'No document text was extracted') + if (Buffer.byteLength(parsed.content, 'utf8') > MAX_DRIVE_TEXT_BYTES) + throw new FileParserError('complexity_limit', 'Extracted text exceeds the live-read limit') + return parsed.content + } catch (error) { + signal?.throwIfAborted() + switch (getFileParserErrorCode(error)) { + case 'encrypted_file': + throw new NativeSearchError( + 'unavailable', + 'This file is password-protected. Remove the password from a copy or open the source.' + ) + case 'empty_input': + case 'no_extractable_text': + throw new NativeSearchError( + 'unavailable', + 'No text could be extracted from this file. Image-only scans require OCR; open the source or provide a searchable copy.' + ) + case 'complexity_limit': + throw new NativeSearchError( + 'unavailable', + 'This file exceeds live text extraction limits. Split it into smaller files or open the source.' + ) + case 'invalid_format': + case 'unsupported_type': + throw new NativeSearchError( + 'unavailable', + 'The file content is malformed or does not match its PDF or DOCX type. Re-save it in a supported format or open the source.' + ) + default: + throw new NativeSearchError( + 'unavailable', + 'The file could not be read as text. Retry or open the source.' + ) + } + } +} diff --git a/apps/sim/lib/sim-search/live/github-service.test.ts b/apps/sim/lib/sim-search/live/github-service.test.ts index 1f0af90787f..908270ee8a7 100644 --- a/apps/sim/lib/sim-search/live/github-service.test.ts +++ b/apps/sim/lib/sim-search/live/github-service.test.ts @@ -39,8 +39,8 @@ const source = { }, } const repository = { id: 123, full_name: 'acme/project', owner: { id: 99 }, default_branch: 'main' } -const member: NativeClient = { json: vi.fn(), text: vi.fn() } -const app: NativeClient = { json: vi.fn(), text: vi.fn() } +const member: NativeClient = { json: vi.fn(), text: vi.fn(), bytes: vi.fn() } +const app: NativeClient = { json: vi.fn(), text: vi.fn(), bytes: vi.fn() } const signal = new AbortController().signal beforeEach(() => { diff --git a/apps/sim/lib/sim-search/live/gitlab-admin.test.ts b/apps/sim/lib/sim-search/live/gitlab-admin.test.ts index e0d719a82ad..5cbcb7d9674 100644 --- a/apps/sim/lib/sim-search/live/gitlab-admin.test.ts +++ b/apps/sim/lib/sim-search/live/gitlab-admin.test.ts @@ -48,7 +48,7 @@ const grant = groupToken({ tenantId: 'gitlab.example.com%3A8443/42', groupId: 'repository', })! -const client = { json: vi.fn(), text: vi.fn() } +const client = { json: vi.fn(), text: vi.fn(), bytes: vi.fn() } const session = (config = source.config) => createAdminGitLabSession({ owner: { organizationId: 'org' }, diff --git a/apps/sim/lib/sim-search/live/gmail-reads.test.ts b/apps/sim/lib/sim-search/live/gmail-reads.test.ts new file mode 100644 index 00000000000..57e955f464f --- /dev/null +++ b/apps/sim/lib/sim-search/live/gmail-reads.test.ts @@ -0,0 +1,142 @@ +import { describe, expect, it } from 'vitest' +import { readGmail } from '@/lib/sim-search/live/google' +import { NativeSearchError } from '@/lib/sim-search/live/http' +import { defaultLiveSearchPolicy } from '@/lib/sim-search/live/policy-schema' +import type { NativeClient } from '@/lib/sim-search/live/types' + +function message(id: string, minute: number, labels = ['INBOX']) { + return { + id, + threadId: 'conversation', + labelIds: labels, + internalDate: String(Date.UTC(2026, 8, 26, 12, minute)), + snippet: `Preview ${id}`, + payload: { + mimeType: 'text/plain', + headers: [ + { name: 'Subject', value: 'Planning discussion' }, + { name: 'From', value: `${id} <${id}@example.test>` }, + ], + body: { data: Buffer.from(`Evidence from ${id}.`).toString('base64url') }, + }, + } +} + +type Message = ReturnType + +function mailbox( + messages: Message[], + options: { + thread?: { id: string; messages: Message[] } + fullMessages?: Record + bodyLimit?: number + } = {} +): NativeClient { + let bodies = 0 + const metadata = ({ payload, ...row }: Message) => ({ + ...row, + payload: { headers: payload.headers }, + }) + return { + async json(path, request) { + if (path === '/gmail/v1/users/me/threads/conversation') { + if (request?.query?.format !== 'metadata') + throw new Error('A thread body request would read messages before authorization') + const thread = options.thread ?? { id: 'conversation', messages } + return { ...thread, messages: thread.messages.map(metadata) } + } + const row = messages.find((entry) => path === `/gmail/v1/users/me/messages/${entry.id}`) + if (!row) throw new Error(`Unexpected Gmail resource: ${path}`) + if (request?.query?.format !== 'full') return metadata(row) + if (++bodies > (options.bodyLimit ?? 8)) throw new Error('Conversation body limit exceeded') + const full = options.fullMessages?.[row.id] ?? row + if (full instanceof Error) throw full + return full + }, + async bytes() { + throw new Error('Unexpected binary request') + }, + async text() { + throw new Error('Gmail messages use the JSON wire format') + }, + } +} + +const readOptions = (verify = async (_reference: { id: string }) => true) => ({ + policy: defaultLiveSearchPolicy('gmail'), + signal: new AbortController().signal, + verify, +}) + +describe('Gmail conversation reads', () => { + it('includes authorized replies with provenance without reading a denied sibling body', async () => { + const client = mailbox([message('reply', 0), message('denied', 2), message('anchor', 1)], { + fullMessages: { denied: new Error('Denied message body was fetched') }, + }) + const result = await readGmail( + client, + 'anchor', + readOptions(async ({ id }) => id !== 'denied') + ) + + expect(result.id).toBe('anchor') + expect(result.accessDependencies).toEqual([{ id: 'reply' }]) + expect(result.content).toContain('Evidence from anchor.') + expect(result.content).toContain('Evidence from reply.') + expect(result.content).not.toContain('Evidence from denied.') + expect(result.content).not.toContain('denied@example.test') + expect(result.content).toContain('reply ') + expect(result.content).toContain('2026-09-26T12:00:00.000Z') + expect(result.content).toContain('https://mail.google.com/mail/u/0/#all/reply') + expect(result.content.indexOf('Evidence from reply.')).toBeLessThan( + result.content.indexOf('Evidence from anchor.') + ) + }) + + it.each(['thread identity', 'missing anchor', 'message identity', 'message membership'] as const)( + 'rejects inconsistent provider evidence: %s', + async (fault) => { + const anchor = message('anchor', 0) + const reply = message('reply', 1) + const client = mailbox([anchor, reply], { + ...(fault === 'thread identity' + ? { thread: { id: 'other-conversation', messages: [anchor, reply] } } + : fault === 'missing anchor' + ? { thread: { id: 'conversation', messages: [reply] } } + : {}), + ...(fault === 'message identity' + ? { fullMessages: { reply: { ...reply, id: 'unrelated' } } } + : fault === 'message membership' + ? { fullMessages: { reply: { ...reply, threadId: 'other-conversation' } } } + : {}), + }) + + await expect(readGmail(client, 'anchor', readOptions())).rejects.toThrow() + } + ) + + it('keeps the requested message when conversation context exceeds the body budget', async () => { + const messages = Array.from({ length: 12 }, (_, index) => message(`mail${index}`, index)) + const result = await readGmail(mailbox(messages), 'mail11', readOptions()) + const included = messages.filter(({ id }) => result.content.includes(`Evidence from ${id}.`)) + + expect(result.id).toBe('mail11') + expect(result.content).toContain('Evidence from mail11.') + expect(included).toHaveLength(8) + expect(result.content).toMatch(/(?:incomplete|limited|omitted|truncated)/i) + }) + + it.each([ + new NativeSearchError('timeout', 'Conversation request timed out'), + new NativeSearchError('rate_limited', 'Gmail quota reached', 45), + ])( + 'preserves provider failure instead of returning a complete conversation: $status', + async (failure) => { + const client = mailbox([message('anchor', 0), message('reply', 1)], { + fullMessages: { reply: failure }, + }) + + await expect(readGmail(client, 'anchor', readOptions())).rejects.toBe(failure) + } + ) +}) diff --git a/apps/sim/lib/sim-search/live/google-service.test.ts b/apps/sim/lib/sim-search/live/google-service.test.ts index 886c18c77aa..b18b6336f88 100644 --- a/apps/sim/lib/sim-search/live/google-service.test.ts +++ b/apps/sim/lib/sim-search/live/google-service.test.ts @@ -15,8 +15,8 @@ vi.mock('@/lib/sim-search/live/http', async (importOriginal) => ({ createNativeClient: mocks.createClient, })) -const member = { json: vi.fn(), text: vi.fn() } -const delegated = { json: vi.fn(), text: vi.fn() } +const member = { json: vi.fn(), text: vi.fn(), bytes: vi.fn() } +const delegated = { json: vi.fn(), text: vi.fn(), bytes: vi.fn() } const mint = vi.fn>() const token = { accessToken: 'directory-token', getDelegatedAccessToken: mint } const person = (email: string) => ({ id: email, email, customerId: 'customer', active: true }) @@ -178,8 +178,13 @@ describe('Google service source filtering', () => { const unavailable = { json: vi.fn().mockRejectedValue(new NativeSearchError('unavailable', 'Source timed out')), text: vi.fn(), + bytes: vi.fn(), + } + const denied = { + json: vi.fn().mockResolvedValue({ id: 'file', trashed: true }), + text: vi.fn(), + bytes: vi.fn(), } - const denied = { json: vi.fn().mockResolvedValue({ id: 'file', trashed: true }), text: vi.fn() } mocks.createClient.mockReturnValueOnce(unavailable).mockReturnValueOnce(denied) const session = await create('google_drive', { userEmails: ['first@example.com', 'second@example.com'], diff --git a/apps/sim/lib/sim-search/live/google.ts b/apps/sim/lib/sim-search/live/google.ts index 9e2c4bf4920..401caf2042c 100644 --- a/apps/sim/lib/sim-search/live/google.ts +++ b/apps/sim/lib/sim-search/live/google.ts @@ -1,3 +1,5 @@ +import { toArray } from '@sim/utils/object' +import { compareStrings, truncate } from '@sim/utils/string' import { mapWithConcurrency } from '@/lib/core/utils/concurrency' import { zonedWallClockToUtc } from '@/lib/core/utils/timezone' import { @@ -7,6 +9,7 @@ import { nativeText, } from '@/lib/sim-search/live/dates' import { readDiscussionSection } from '@/lib/sim-search/live/discussion' +import { readDriveFileContent } from '@/lib/sim-search/live/drive-content' import { array, NativeSearchError, object, segment, string } from '@/lib/sim-search/live/http' import { interleaveByRank } from '@/lib/sim-search/live/pages' import { permitsResources } from '@/lib/sim-search/live/policy' @@ -15,6 +18,7 @@ import type { NativeClient, NativeDocument, NativePage, + NativeReadOptions, NativeSearchInput, } from '@/lib/sim-search/live/types' @@ -81,10 +85,17 @@ export async function searchDrive( } } -export async function readDrive(client: NativeClient, id: string): Promise { +export async function readDrive( + client: NativeClient, + id: string, + signal?: AbortSignal +): Promise { const row = object( await client.json(`/drive/v3/files/${segment(id)}`, { - query: { fields: DRIVE_FIELDS, supportsAllDrives: 'true' }, + query: { + fields: `${DRIVE_FIELDS},size,capabilities(canDownload)`, + supportsAllDrives: 'true', + }, }) ) const document = driveDocument(row) @@ -117,6 +128,11 @@ export async function readDrive(client: NativeClient, id: string): Promise `${string(range.range)}\n${JSON.stringify(range.values ?? [])}`) .join('\n\n') + } else if ( + document.kind === 'application/pdf' || + document.kind === 'application/vnd.openxmlformats-officedocument.wordprocessingml.document' + ) { + document.content = await readDriveFileContent(client, id, row, signal) } else if (document.kind?.startsWith('text/') || document.kind === 'application/json') { document.content = await client.text(`/drive/v3/files/${segment(id)}`, { alt: 'media', @@ -265,25 +281,201 @@ export async function searchGmail( } } -function mailText(value: unknown): string { - const part = object(value) - const children = array(part.parts).map(mailText).filter(Boolean) - const encoded = string(object(part.body).data) - if (string(part.mimeType) === 'text/plain' && encoded) - return providerText(Buffer.from(encoded, 'base64url').toString('utf8')) - if (children.length) return children.join('\n') - // HTML-only messages remain text, never rendered as markup. - if (string(part.mimeType) === 'text/html' && encoded) - return providerText(Buffer.from(encoded, 'base64url').toString('utf8'), 'html') - return '' +/** The response byte cap does not bound recursive MIME depth or per-part processing work. */ +const GMAIL_MIME_DEPTH_LIMIT = 32 +const GMAIL_MIME_PART_LIMIT = 256 + +interface GmailBodyText { + text: string + incomplete: boolean + plain: boolean } -export async function readGmail(client: NativeClient, id: string): Promise { - const data = object( +/** Selects one available alternative, preserves mixed body sections, and omits file attachments. */ +function mailText(value: unknown): GmailBodyText { + let remainingParts = GMAIL_MIME_PART_LIMIT + const bounded = (text: string, incomplete: boolean, plain: boolean): GmailBodyText => ({ + text: truncate(text, GMAIL_MESSAGE_CHARACTER_LIMIT), + incomplete: incomplete || text.length > GMAIL_MESSAGE_CHARACTER_LIMIT, + plain, + }) + const visit = (value: unknown, depth: number): GmailBodyText => { + if (remainingParts === 0) return { text: '', incomplete: true, plain: false } + remainingParts-- + if (depth >= GMAIL_MIME_DEPTH_LIMIT) return { text: '', incomplete: true, plain: false } + const part = object(value) + if (string(part.filename)) return { text: '', incomplete: false, plain: false } + const mimeType = string(part.mimeType).toLowerCase() + if (mimeType === 'text/plain' || mimeType === 'text/html') { + const body = object(part.body) + const encoded = string(body.data) + return bounded( + encoded + ? providerText( + Buffer.from(encoded, 'base64url').toString('utf8'), + mimeType === 'text/plain' ? 'plain' : 'html' + ) + : '', + !encoded && (Boolean(body.attachmentId) || Number(body.size) > 0), + mimeType === 'text/plain' + ) + } + const children: GmailBodyText[] = [] + let incomplete = false + for (const child of toArray(part.parts)) { + if (remainingParts === 0) { + incomplete = true + break + } + children.push(visit(child, depth + 1)) + } + if (mimeType === 'multipart/alternative') { + const selected = + children.find((child) => child.text && child.plain) ?? children.find((child) => child.text) + if (selected) return { ...selected, incomplete: selected.incomplete || incomplete } + } + return bounded( + children + .map((child) => child.text) + .filter(Boolean) + .join('\n'), + incomplete || children.some((child) => child.incomplete), + children.some((child) => child.text && child.plain) + ) + } + return visit(value, 0) +} + +/** Eight messages leave room for fresh per-message scope checks in the member request budget. */ +const GMAIL_CONVERSATION_MESSAGE_LIMIT = 8 +/** Per-message bounds retain the anchor even when earlier messages contain large bodies. */ +const GMAIL_MESSAGE_CHARACTER_LIMIT = 24_000 + +/** + * The application has authorized the anchor. Thread discovery returns metadata only; each + * sibling is authorized before its body is fetched and retained for the final fresh check. + */ +export async function readGmail( + client: NativeClient, + id: string, + options: NativeReadOptions +): Promise { + options.signal.throwIfAborted() + const anchor = object( await client.json(`/gmail/v1/users/me/messages/${segment(id)}`, { query: { format: 'full' } }) ) - const document = gmailDocument(data) - return { ...document, content: mailText(data.payload) || document.content } + const threadId = string(anchor.threadId) + if ( + anchor.id !== id || + typeof anchor.threadId !== 'string' || + !threadId || + threadId.length > 1000 + ) + throw new NativeSearchError('unavailable', 'Gmail returned inconsistent message identity.') + const thread = object( + await client.json(`/gmail/v1/users/me/threads/${segment(threadId)}`, { + query: { + format: 'metadata', + fields: 'id,messages(id,threadId,labelIds,internalDate)', + }, + }) + ) + if (thread.id !== threadId || !Array.isArray(thread.messages)) + throw new NativeSearchError('unavailable', 'Gmail returned inconsistent conversation identity.') + const members = new Map>() + for (const value of thread.messages) { + const row = object(value) + if (typeof row.id !== 'string' || !row.id || row.id.length > 1000 || row.threadId !== threadId) + throw new NativeSearchError( + 'unavailable', + 'Gmail returned inconsistent conversation membership.' + ) + if (!members.has(row.id)) members.set(row.id, row) + } + if (!members.has(id)) + throw new NativeSearchError( + 'unavailable', + 'The requested message is missing from its conversation.' + ) + const latest = [...members.values()] + .filter((row) => row.id !== id) + .sort( + (left, right) => + (Number(right.internalDate) || 0) - (Number(left.internalDate) || 0) || + compareStrings(string(left.id), string(right.id)) + ) + .slice(0, GMAIL_CONVERSATION_MESSAGE_LIMIT - 1) + const warnings = new Set() + if (members.size > GMAIL_CONVERSATION_MESSAGE_LIMIT) + warnings.add('only the requested message and up to seven recent messages are included') + const documents: NativeDocument[] = [] + const append = (row: Record) => { + const document = gmailDocument(row) + const body = mailText(row.payload) + if (body.incomplete) warnings.add('some message body content could not be read') + if (!body.text && document.content) warnings.add('some messages contain previews only') + const content = [ + `## Message ${document.id}`, + `From: ${truncate(document.author || 'Unknown sender', 1000)}`, + `Date: ${document.modifiedAt || 'Unknown date'}`, + `Subject: ${truncate(document.title, 1000)}`, + `Source: ${document.url}`, + '', + body.text || + (document.content + ? `Preview only: ${document.content}` + : 'No inline message body available.'), + ].join('\n') + if (content.length > GMAIL_MESSAGE_CHARACTER_LIMIT) warnings.add('message text was truncated') + document.content = truncate(content, GMAIL_MESSAGE_CHARACTER_LIMIT) + documents.push(document) + return document + } + const document = append(anchor) + for (const candidate of latest) { + options.signal.throwIfAborted() + const candidateId = string(candidate.id) + if ( + !(await options.verify({ + id: candidateId, + accessMetadata: { id: candidateId, labelIds: candidate.labelIds }, + })) + ) { + warnings.add('messages outside the source search scope were omitted') + continue + } + const row = object( + await client.json(`/gmail/v1/users/me/messages/${segment(candidateId)}`, { + query: { format: 'full' }, + }) + ) + if (row.id !== candidateId || row.threadId !== threadId) + throw new NativeSearchError( + 'unavailable', + 'Gmail returned inconsistent conversation membership.' + ) + append(row) + } + options.signal.throwIfAborted() + documents.sort( + (left, right) => + (Date.parse(left.modifiedAt ?? '') || 0) - (Date.parse(right.modifiedAt ?? '') || 0) || + compareStrings(left.id, right.id) + ) + return { + ...document, + accessDependencies: documents + .filter((entry) => entry.id !== id) + .map((entry) => ({ id: entry.id })), + content: [ + ...(warnings.size + ? [ + `Coverage incomplete: ${[...warnings].join('; ')}. Open Gmail for the remaining context.`, + ] + : []), + ...documents.map((entry) => entry.content), + ].join('\n\n'), + } } function eventDocument( diff --git a/apps/sim/lib/sim-search/live/http.ts b/apps/sim/lib/sim-search/live/http.ts index f53d71aa086..81f20d50654 100644 --- a/apps/sim/lib/sim-search/live/http.ts +++ b/apps/sim/lib/sim-search/live/http.ts @@ -19,6 +19,8 @@ export class NativeSearchError extends Error { /** Provider requests one native search may make, including discovery and verification. */ export const NATIVE_SEARCH_REQUEST_BUDGET = 30 +/** Decoded bytes allowed in any native provider response, including file downloads. */ +export const NATIVE_RESPONSE_MAX_BYTES = 4 * 1024 * 1024 /** Tokens only go to a code-selected provider origin; redirects never carry credentials. */ export function createNativeClient(input: { @@ -68,7 +70,7 @@ export function createNativeClient(input: { ...(options?.body ? { body: JSON.stringify(options.body) } : {}), signal: input.signal, timeout: 10_000, - maxResponseBytes: 4 * 1024 * 1024, + maxResponseBytes: NATIVE_RESPONSE_MAX_BYTES, maxRedirects: 0, acceptCompressed: true, connectionPool: input.pool, @@ -117,6 +119,14 @@ export function createNativeClient(input: { async json(path, options) { return (await request(path, options)).json() }, + async bytes(path, query) { + try { + return Buffer.from(await (await request(path, { query })).arrayBuffer()) + } catch (error) { + input.signal.throwIfAborted() + throw error + } + }, async text(path, query) { return (await request(path, { query })).text() }, diff --git a/apps/sim/lib/sim-search/live/linear.test.ts b/apps/sim/lib/sim-search/live/linear.test.ts index 5c472da4430..a8536527865 100644 --- a/apps/sim/lib/sim-search/live/linear.test.ts +++ b/apps/sim/lib/sim-search/live/linear.test.ts @@ -18,6 +18,9 @@ const input: NativeSearchInput = { query: 'rollout', limit: 10, scopes: [] } function client(json: NativeClient['json']): NativeClient { return { json, + bytes: async () => { + throw new Error('Unexpected binary request') + }, text: async () => { throw new Error('Unexpected text request') }, diff --git a/apps/sim/lib/sim-search/live/policy.test.ts b/apps/sim/lib/sim-search/live/policy.test.ts index 491c3edecc9..5d350fa0898 100644 --- a/apps/sim/lib/sim-search/live/policy.test.ts +++ b/apps/sim/lib/sim-search/live/policy.test.ts @@ -10,6 +10,7 @@ function client(rows: Record): NativeClient { return rows[path] }), text: vi.fn(), + bytes: vi.fn(), } } const selected = (included: string[], excluded: string[] = []) => ({ diff --git a/apps/sim/lib/sim-search/live/providers.test.ts b/apps/sim/lib/sim-search/live/providers.test.ts index 2e0df1529aa..a59aeece2c2 100644 --- a/apps/sim/lib/sim-search/live/providers.test.ts +++ b/apps/sim/lib/sim-search/live/providers.test.ts @@ -5,7 +5,11 @@ import { searchSlack } from '@/lib/sim-search/live/slack' import type { NativeClient } from '@/lib/sim-search/live/types' function client() { - return { json: vi.fn(), text: vi.fn() } + return { + json: vi.fn(), + text: vi.fn(), + bytes: vi.fn(), + } } const input = { query: 'launch', limit: 20, scopes: [] } diff --git a/apps/sim/lib/sim-search/live/providers.ts b/apps/sim/lib/sim-search/live/providers.ts index 3c4f8463a42..ffab433e4c7 100644 --- a/apps/sim/lib/sim-search/live/providers.ts +++ b/apps/sim/lib/sim-search/live/providers.ts @@ -1,4 +1,3 @@ -import type { WorkspaceSearchFilters } from '@/lib/api/contracts/knowledge' import { MAX_NATIVE_QUERIES_PER_ACCOUNT } from '@/lib/api/contracts/mothership-assistant-tools' import { readAtlassian, searchAtlassian } from '@/lib/sim-search/live/atlassian' import { readCoda, searchCoda } from '@/lib/sim-search/live/coda' @@ -14,7 +13,6 @@ import { } from '@/lib/sim-search/live/google' import { NativeSearchError } from '@/lib/sim-search/live/http' import { readLinear, searchLinear } from '@/lib/sim-search/live/linear' -import type { LiveSearchPolicy } from '@/lib/sim-search/live/policy-schema' import { LIVE_SEARCH_PROVIDER_IDS, type LiveSearchProviderId, @@ -24,6 +22,7 @@ import type { NativeClient, NativeDocument, NativePage, + NativeReadOptions, NativeSearchInput, } from '@/lib/sim-search/live/types' @@ -49,8 +48,7 @@ interface NativeProvider { read( client: NativeClient, reference: Pick, - policy?: LiveSearchPolicy, - filters?: WorkspaceSearchFilters + options: NativeReadOptions ): Promise } @@ -69,10 +67,10 @@ export const LIVE_SEARCH_PROVIDERS = { "'person@example.com' in owners (or writers, readers), mimeType = 'application/vnd.google-apps.document' (or spreadsheet, presentation, folder) and 'FOLDER_ID' in parents; project drive:DRIVE_ID searches one shared drive, whose files have no owners.", example: "fullText contains 'roadmap' and 'jane@example.com' in owners", avoid: - 'bare words without a term and operator, which Drive rejects, and trashed or modifiedTime clauses, which the server adds from startDate/endDate. Drive search does not search comments or replies; find the file by title/content, then read it to retrieve its discussion.', + 'bare words without a term and operator, which Drive rejects, and trashed or modifiedTime clauses, which the server adds from startDate/endDate. Drive search does not search comments or replies; find the file by title/content, then read it to retrieve its discussion. PDF and DOCX reads extract text within download and parsing limits; scanned PDFs need OCR and unsupported binaries provide metadata only.', }, search: searchDrive, - read: (client, reference) => readDrive(client, reference.id), + read: (client, reference, options) => readDrive(client, reference.id, options.signal), }, gmail: { guide: { @@ -82,10 +80,10 @@ export const LIVE_SEARCH_PROVIDERS = { 'from:, to:, cc:, subject:, label:, has:attachment, filename:, in:sent, from:me and is:unread; spam and trash are not searched.', example: 'from:jane@example.com subject:(budget OR forecast)', avoid: - 'listing alternatives with spaces, which requires all of them, or lowercase or; join alternatives with uppercase OR.', + 'listing alternatives with spaces, which requires all of them, or lowercase or; join alternatives with uppercase OR. Search matches individual messages. Reading a match includes up to eight independently authorized messages from its conversation, with per-message dates and URLs. Context can fall outside the search dates or keywords; use the message timestamps and incomplete notices, and never claim the entire thread was searched.', }, search: searchGmail, - read: (client, reference) => readGmail(client, reference.id), + read: (client, reference, options) => readGmail(client, reference.id, options), }, google_calendar: { guide: { @@ -98,13 +96,13 @@ export const LIVE_SEARCH_PROVIDERS = { 'OR, quotes or field operators, which q does not support; run alternatives as separate native queries.', }, search: searchCalendar, - read: (client, reference, policy, filters) => + read: (client, reference, options) => readCalendar( client, reference.id, reference.container, - policy?.includeAttendees, - Boolean(filters?.startDate || filters?.endDate) + options.policy.includeAttendees, + Boolean(options.filters?.startDate || options.filters?.endDate) ), }, slack: { @@ -257,8 +255,7 @@ export function readNativeProvider( provider: LiveSearchProviderId, client: NativeClient, reference: Pick, - policy?: LiveSearchPolicy, - filters?: WorkspaceSearchFilters + options: NativeReadOptions ): Promise { const adapter = LIVE_SEARCH_PROVIDERS[provider] if (!('read' in adapter)) @@ -266,7 +263,7 @@ export function readNativeProvider( 'unavailable', `${provider} requires a connected member MCP account` ) - return adapter.read(client, reference, policy, filters) + return adapter.read(client, reference, options) } /** Rules for every provider, ahead of the query cards of the providers in play. */ diff --git a/apps/sim/lib/sim-search/live/service-session.test.ts b/apps/sim/lib/sim-search/live/service-session.test.ts index 6ee659a07e9..fe361a672b8 100644 --- a/apps/sim/lib/sim-search/live/service-session.test.ts +++ b/apps/sim/lib/sim-search/live/service-session.test.ts @@ -41,7 +41,7 @@ vi.mock('@/connectors/registry', () => ({ }, })) -const api = { json: vi.fn(), text: vi.fn() } +const api = { json: vi.fn(), text: vi.fn(), bytes: vi.fn() } const source = { id: 'source', credentialId: 'service', diff --git a/apps/sim/lib/sim-search/live/slack-format.test.ts b/apps/sim/lib/sim-search/live/slack-format.test.ts index 44587773a31..c4222175eb2 100644 --- a/apps/sim/lib/sim-search/live/slack-format.test.ts +++ b/apps/sim/lib/sim-search/live/slack-format.test.ts @@ -36,6 +36,7 @@ describe('Slack search presentation', () => { throw aborted }), text: vi.fn(), + bytes: vi.fn(), } await expect(readSlack(api, '123.456', 'D1')).rejects.toBe(aborted) }) diff --git a/apps/sim/lib/sim-search/live/types.ts b/apps/sim/lib/sim-search/live/types.ts index f146546bf6b..54afa4a6704 100644 --- a/apps/sim/lib/sim-search/live/types.ts +++ b/apps/sim/lib/sim-search/live/types.ts @@ -20,7 +20,20 @@ export interface LiveAccount { scopes: string[] } +export type NativeAccessReference = Pick + +export interface NativeReadOptions { + policy: LiveSearchPolicy + filters?: WorkspaceSearchFilters + signal: AbortSignal + verify( + reference: NativeAccessReference & Pick + ): Promise +} + export interface NativeDocument { + /** Server-only references whose content must pass a fresh check before projection. */ + accessDependencies?: readonly NativeAccessReference[] /** * Server-only permission evidence from the same response as the content. Only a verifier * bound to the client that produced it may consume it; it is never projected to callers. @@ -69,6 +82,7 @@ export interface NativeClient { } ): Promise text(path: string, query?: Record): Promise + bytes(path: string, query?: Record): Promise } export interface NativeSearchInput { diff --git a/apps/sim/scripts/test-search-discussions-live.ts b/apps/sim/scripts/test-search-discussions-live.ts index 310c1372875..440a414cec0 100644 --- a/apps/sim/scripts/test-search-discussions-live.ts +++ b/apps/sim/scripts/test-search-discussions-live.ts @@ -179,6 +179,9 @@ const client: NativeClient = { const endpoint = `${path}${query.size ? `?${query}` : ''}` return githubApi(endpoint, 'adapter') }, + async bytes() { + throw new Error('Unexpected binary request in a PR read') + }, async text() { throw new Error('Unexpected text request in a PR read') }, diff --git a/apps/sim/scripts/test-search-google-content-e2e.ts b/apps/sim/scripts/test-search-google-content-e2e.ts new file mode 100644 index 00000000000..ad968f3f7d2 --- /dev/null +++ b/apps/sim/scripts/test-search-google-content-e2e.ts @@ -0,0 +1,804 @@ +import assert from 'node:assert/strict' +import { randomBytes } from 'node:crypto' +import { mkdir, writeFile } from 'node:fs/promises' +import http from 'node:http' +import type { AddressInfo } from 'node:net' +import { dirname } from 'node:path' +import { createLogger } from '@sim/logger' +import { createDeferred } from '@sim/testing/helpers/deferred' +import { getErrorMessage } from '@sim/utils/errors' +import { truncate } from '@sim/utils/string' +import JSZip from 'jszip' +import { PDFDocument, StandardFonts } from 'pdf-lib' +import { isHosted } from '@/lib/core/config/env-flags' +import { DocxParser } from '@/lib/file-parsers/docx-parser' +import { PdfParser } from '@/lib/file-parsers/pdf-parser' +import { readDrive, readGmail } from '@/lib/sim-search/live/google' +import { createNativeClient, NativeSearchError } from '@/lib/sim-search/live/http' +import { createPolicyVerifier } from '@/lib/sim-search/live/policy' +import { defaultLiveSearchPolicy } from '@/lib/sim-search/live/policy-schema' + +/** + * Exercises Drive content reads over a real loopback HTTP server with real document parsers. + * Run from apps/sim: + * NEXT_PUBLIC_APP_URL=http://localhost:3040 NEXT_PUBLIC_FORCE_HOSTED=false \ + * SEARCH_GOOGLE_CONTENT_REPORT_PATH=/tmp/google-content.json \ + * bun scripts/test-search-google-content-e2e.ts + * Fixtures are synthetic. No Google credentials or remote writes are used. This verifies the + * native transport, Drive adapter and parsers; connected-account authorization and UI are + * separate acceptance boundaries. The server uses the existing self-hosted loopback policy. + */ +const logger = createLogger('SearchGoogleContentE2E') +const reportPath = process.env.SEARCH_GOOGLE_CONTENT_REPORT_PATH +assert(reportPath, 'Set SEARCH_GOOGLE_CONTENT_REPORT_PATH') +assert(!isHosted, 'Use a local self-hosted app URL with NEXT_PUBLIC_FORCE_HOSTED=false') +const PDF_MIME = 'application/pdf' +const DOCX_MIME = 'application/vnd.openxmlformats-officedocument.wordprocessingml.document' +const MIB = 1024 * 1024 +const TOKEN = 'synthetic-loopback-token' +const SOURCE = 'https://drive.google.com/file/d/synthetic-file/view' +const COMMENT = 'Synthetic reviewer requests a revised launch estimate.' +const checks: { name: string; status: 'passed' | 'failed'; durationMs: number; error?: string }[] = + [] +const requests: { + id: string + resource: 'metadata' | 'media' | 'comments' | 'message' | 'thread' | 'labels' | 'attachment' + status: number +}[] = [] + +interface Fixture { + mimeType: string + body: Buffer + metadata?: Record + mode?: 'stall' | 'overflow' +} +const fixtures = new Map() +interface GmailPart { + mimeType: string + filename?: string + headers?: { name: string; value: string }[] + body?: { data?: string; attachmentId?: string; size?: number } + parts?: GmailPart[] +} +interface GmailMessage { + id: string + threadId: string + labelIds: string[] + internalDate: string + snippet?: string + payload: GmailPart +} +const mailbox: GmailMessage[] = Array.from({ length: 12 }, (_, index) => ({ + id: `mail${index}`, + threadId: 'conversation', + labelIds: ['INBOX'], + internalDate: String(Date.UTC(2026, 0, 1, 12, index)), + payload: { + mimeType: 'text/plain', + headers: [{ name: 'Subject', value: 'Synthetic planning' }], + body: { data: Buffer.from(`Synthetic evidence ${index}.`).toString('base64url') }, + }, +})) +const mediaStarted = createDeferred() +const mediaClosed = createDeferred() +const server = http.createServer((request, response) => { + if (request.headers.authorization !== `Bearer ${TOKEN}`) { + response.writeHead(401).end() + return + } + const url = new URL(request.url ?? '/', 'http://127.0.0.1') + if (url.pathname.startsWith('/gmail/v1/users/me/')) { + const id = url.pathname.split('/').at(-1) ?? '' + const row = mailbox.find((message) => message.id === id) + const isThread = url.pathname.startsWith('/gmail/v1/users/me/threads/') + const isAttachment = url.pathname.includes('/attachments/') + const threadMessages = isThread ? mailbox.filter((message) => message.threadId === id) : [] + const metadata = (message: (typeof mailbox)[number]) => ({ + ...message, + payload: { headers: message.payload.headers }, + }) + const data = + id === 'labels' + ? { labels: [{ id: 'INBOX', name: 'INBOX' }] } + : isThread && threadMessages.length > 0 && url.searchParams.get('format') === 'metadata' + ? { id, messages: threadMessages.map(metadata) } + : row + ? url.searchParams.get('format') === 'full' + ? row + : metadata(row) + : undefined + requests.push({ + id, + resource: isAttachment + ? 'attachment' + : id === 'labels' + ? 'labels' + : isThread + ? 'thread' + : 'message', + status: data ? 200 : 400, + }) + response + .writeHead(data ? 200 : 400, { 'Content-Type': 'application/json' }) + .end(JSON.stringify(data ?? {})) + return + } + const match = /^\/drive\/v3\/files\/([^/]+)(\/comments)?$/.exec(url.pathname) + const id = match?.[1] ?? '' + const fixture = fixtures.get(id) + if (!fixture) { + response.writeHead(404).end() + return + } + const resource = match?.[2] + ? 'comments' + : url.searchParams.get('alt') === 'media' + ? 'media' + : 'metadata' + requests.push({ id, resource, status: 200 }) + if (resource === 'media') { + response.writeHead(200, { 'Content-Type': fixture.mimeType }) + if (fixture.mode === 'stall') { + request.socket.once('close', () => mediaClosed.resolve()) + response.write(fixture.body.subarray(0, 8)) + mediaStarted.resolve() + } else if (fixture.mode === 'overflow') { + response.write(Buffer.alloc(4 * MIB + 1, 0x41)) + } else response.end(fixture.body) + return + } + response.setHeader('Content-Type', 'application/json') + const fields = url.searchParams.get('fields') ?? '' + const metadata = { + id, + name: 'Synthetic document title', + description: 'Synthetic metadata description only.', + mimeType: fixture.mimeType, + size: String(fixture.body.length), + capabilities: { canDownload: true }, + webViewLink: SOURCE, + modifiedTime: '2026-01-01T12:00:00Z', + ...fixture.metadata, + } + if (!fields.includes('size') && fields !== '*') Reflect.deleteProperty(metadata, 'size') + if (!fields.includes('capabilities') && fields !== '*') + Reflect.deleteProperty(metadata, 'capabilities') + response.end( + JSON.stringify( + resource === 'comments' + ? { + comments: [ + { + id: 'synthetic-comment', + content: COMMENT, + createdTime: '2026-01-02T12:00:00Z', + author: { displayName: 'Synthetic Reviewer' }, + }, + ], + } + : metadata + ) + ) +}) + +async function check(name: string, run: () => Promise) { + const started = performance.now() + try { + await run() + checks.push({ name, status: 'passed', durationMs: Math.round(performance.now() - started) }) + } catch (error) { + checks.push({ + name, + status: 'failed', + durationMs: Math.round(performance.now() - started), + error: truncate(getErrorMessage(error), 1000), + }) + } +} + +async function pdf(pages: string[]): Promise { + const document = await PDFDocument.create() + const font = await document.embedFont(StandardFonts.Helvetica) + for (const text of pages) { + const page = document.addPage([612, 792]) + if (text) + page.drawText(text, { x: 40, y: 720, font, size: text.length > 1000 ? 4 : 12, lineHeight: 8 }) + else page.drawRectangle({ x: 40, y: 40, width: 100, height: 100 }) + } + return Buffer.from(await document.save()) +} + +async function docx(body: string, parts: Record = {}): Promise { + const zip = new JSZip() + zip.file( + '[Content_Types].xml', + '' + ) + zip.file( + '_rels/.rels', + '' + ) + zip.file( + 'word/document.xml', + `${body}` + ) + for (const [name, bytes] of Object.entries(parts)) zip.file(name, bytes) + return zip.generateAsync({ type: 'nodebuffer', compression: 'DEFLATE' }) +} + +/** A real password-required PDF structure; no decrypted content is embedded. */ +function encryptedPdf(): Buffer { + const key = '00'.repeat(32) + const objects = [ + '<< /Type /Catalog /Pages 2 0 R >>', + '<< /Type /Pages /Kids [3 0 R] /Count 1 >>', + '<< /Type /Page /Parent 2 0 R /MediaBox [0 0 612 792] >>', + `<< /Filter /Standard /V 1 /R 2 /O <${key}> /U <${key}> /P -4 >>`, + ] + let text = '%PDF-1.4\n' + const offsets = objects.map((object, index) => { + const offset = Buffer.byteLength(text) + text += `${index + 1} 0 obj\n${object}\nendobj\n` + return offset + }) + const xref = Buffer.byteLength(text) + text += `xref\n0 5\n0000000000 65535 f \n${offsets.map((offset) => `${String(offset).padStart(10, '0')} 00000 n \n`).join('')}trailer\n<< /Size 5 /Root 1 0 R /Encrypt 4 0 R /ID [<${'11'.repeat(16)}> <${'11'.repeat(16)}>] >>\nstartxref\n${xref}\n%%EOF\n` + return Buffer.from(text) +} + +function paragraph(text: string): string { + return `${text}` +} + +let origin = '' +async function read(id: string, signal = AbortSignal.timeout(5000)) { + const api = createNativeClient({ origin, accessToken: TOKEN, signal }) + return readDrive(api, id, signal) +} +async function readMailFixture(payload: GmailPart, snippet?: string) { + const id = `mime${mailbox.length}` + mailbox.push({ + id, + threadId: `conversation-${id}`, + labelIds: ['INBOX'], + internalDate: String(Date.UTC(2026, 0, 1, 12, 0)), + snippet, + payload, + }) + const signal = AbortSignal.timeout(5000) + const api = createNativeClient({ origin, accessToken: TOKEN, signal }) + const policy = defaultLiveSearchPolicy('gmail') + const verify = createPolicyVerifier('gmail', policy, api, origin) + return readGmail(api, id, { policy, signal, verify }) +} +async function unavailable(id: string) { + await assert.rejects( + read(id), + (error: unknown) => error instanceof NativeSearchError && error.status === 'unavailable' + ) +} + +try { + const started = createDeferred() + server.once('error', started.reject) + server.listen(0, '127.0.0.1', () => started.resolve()) + await started.promise + origin = `http://127.0.0.1:${(server.address() as AddressInfo).port}` + await check( + 'Gmail conversation and fresh selected-label checks fit one native request budget', + async () => { + const signal = AbortSignal.timeout(5000) + const api = createNativeClient({ origin, accessToken: TOKEN, signal }) + const policy = { + ...defaultLiveSearchPolicy('gmail'), + mode: 'selected' as const, + included: ['INBOX'], + } + const verify = createPolicyVerifier('gmail', policy, api, origin) + assert.ok(await verify({ id: 'mail11' })) + const document = await readGmail(api, 'mail11', { + policy, + signal, + verify: (reference) => verify(reference, reference.accessMetadata), + }) + assert.equal(document.id, 'mail11') + assert.equal( + mailbox.filter((_message, index) => + document.content.includes(`Synthetic evidence ${index}.`) + ).length, + 8 + ) + const current = createPolicyVerifier('gmail', policy, api, origin, undefined, { fresh: true }) + assert.ok(await current({ id: document.id })) + for (const reference of document.accessDependencies ?? []) assert.ok(await current(reference)) + } + ) + await check( + 'Gmail MIME alternatives contribute one body rather than duplicate evidence', + async () => { + const evidence = 'Synthetic MIME alternative evidence.' + const document = await readMailFixture({ + mimeType: 'multipart/alternative', + parts: [ + { + mimeType: 'text/html', + body: { data: Buffer.from(`

${evidence}

`).toString('base64url') }, + }, + { + mimeType: 'text/plain', + body: { data: Buffer.from(evidence).toString('base64url') }, + }, + ], + }) + assert.equal(document.content.split(evidence).length - 1, 1) + } + ) + await check('Gmail named text attachments do not become message evidence', async () => { + const evidence = 'Synthetic message body remains readable.' + const attachment = 'Synthetic attachment content is not a message statement.' + const document = await readMailFixture({ + mimeType: 'multipart/mixed', + parts: [ + { + mimeType: 'text/plain', + body: { data: Buffer.from(evidence).toString('base64url') }, + }, + { + mimeType: 'text/plain', + filename: 'synthetic-notes.txt', + body: { data: Buffer.from(attachment).toString('base64url') }, + }, + ], + }) + assert.ok(document.content.includes(evidence)) + assert.ok(!document.content.includes(attachment)) + }) + await check( + 'Gmail external body previews disclose missing content without attachment fanout', + async () => { + const before = requests.length + const preview = 'Synthetic preview only.' + const document = await readMailFixture( + { mimeType: 'text/plain', body: { attachmentId: 'synthetic-body', size: 1200 } }, + preview + ) + assert.ok(document.content.includes(preview)) + assert.match(document.content, /incomplete/i) + assert.equal( + requests.slice(before).filter((request) => request.resource === 'attachment').length, + 0 + ) + } + ) + await check('Gmail MIME depth limits omit deep content with an incomplete notice', async () => { + const evidence = 'Synthetic deeply nested body must be omitted.' + let payload: GmailPart = { + mimeType: 'text/plain', + body: { data: Buffer.from(evidence).toString('base64url') }, + } + for (let depth = 0; depth < 40; depth++) + payload = { mimeType: 'multipart/mixed', parts: [payload] } + const document = await readMailFixture(payload, 'Synthetic nested body preview.') + assert.match(document.content, /incomplete/i) + assert.ok(!document.content.includes(evidence)) + }) + await check( + 'Gmail MIME part limits preserve early sections and disclose omitted content', + async () => { + const document = await readMailFixture({ + mimeType: 'multipart/mixed', + parts: Array.from({ length: 300 }, (_, index) => ({ + mimeType: 'text/plain', + body: { data: Buffer.from(`Synthetic section ${index}.`).toString('base64url') }, + })), + }) + assert.ok(document.content.includes('Synthetic section 0.')) + assert.ok(!document.content.includes('Synthetic section 299.')) + assert.match(document.content, /incomplete/i) + } + ) + const pdfBody = await pdf([ + 'The synthetic launch budget is forty-two credits.', + 'The launch owner is the synthetic research team.', + ]) + const wordBody = await docx( + paragraph('The synthetic decision is to launch on Tuesday.') + + '' + + paragraph('Synthetic table budget') + + '' + + paragraph('Seventy credits') + + '' + ) + fixtures.set('pdf', { mimeType: PDF_MIME, body: pdfBody }) + fixtures.set('docx', { mimeType: DOCX_MIME, body: wordBody }) + + for (const [id, evidence] of [ + [ + 'pdf', + [ + 'The synthetic launch budget is forty-two credits.', + 'The launch owner is the synthetic research team.', + ], + ], + [ + 'docx', + [ + 'The synthetic decision is to launch on Tuesday.', + 'Synthetic table budget', + 'Seventy credits', + ], + ], + ] as const) { + await check(`${id} bytes produce attributed document text and retain discussion`, async () => { + const document = await read(id) + for (const expected of [...evidence, SOURCE, COMMENT]) + assert.ok(document.content.includes(expected), `Missing synthetic evidence: ${expected}`) + assert.equal(document.id, id) + assert.equal(document.url, SOURCE) + }) + } + + const noteReference = (id: number) => `` + const notesPart = (body: string) => ({ + 'word/footnotes.xml': Buffer.from( + `${body}` + ), + }) + const repeatedNote = 'Synthetic footnote evidence survives each reference.' + fixtures.set('docx-repeated-note', { + mimeType: DOCX_MIME, + body: await docx( + `Synthetic cited body.${noteReference(1)}${noteReference(1)}`, + notesPart(paragraph(repeatedNote)) + ), + }) + await check('DOCX reads retain each referenced note within the extraction budget', async () => { + const document = await read('docx-repeated-note') + assert.ok(document.content.includes('Synthetic cited body.')) + assert.equal(document.content.split(repeatedNote).length - 1, 2) + }) + + const paddedImage = Buffer.concat([ + Buffer.from( + 'iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAAC0lEQVR4nGP4DwQACfsD/fteaysAAAAASUVORK5CYII=', + 'base64' + ), + randomBytes(MIB), + ]) + const imageReferences = Array.from( + { length: 25 }, + (_, index) => + `` + ).join('') + fixtures.set('docx-repeated-image', { + mimeType: DOCX_MIME, + body: await docx( + `Synthetic body with image references.${noteReference(1)}${imageReferences}`, + { + ...notesPart(paragraph(repeatedNote)), + '[Content_Types].xml': Buffer.from( + '' + ), + 'word/_rels/document.xml.rels': Buffer.from( + '' + ), + 'word/media/synthetic.png': paddedImage, + } + ), + }) + await check('DOCX image references cannot displace readable body and note evidence', async () => { + const document = await read('docx-repeated-image') + assert.ok(document.content.includes('Synthetic body with image references.')) + assert.ok(document.content.includes(repeatedNote)) + }) + + const amplifiedNote = randomBytes(225_000).toString('base64') + fixtures.set('docx-amplified-note', { + mimeType: DOCX_MIME, + body: await docx( + `Synthetic cited body.${noteReference(1).repeat(32)}${noteReference(2)}`, + notesPart(paragraph(amplifiedNote)) + ), + }) + fixtures.set('docx-structural-complexity', { + mimeType: DOCX_MIME, + body: await docx( + paragraph('Synthetic body before excessive structural elements.') + + `${''.repeat(50_001)}` + ), + }) + for (const id of ['docx-amplified-note', 'docx-structural-complexity']) { + await check(`${id} rejects before a fallback can return incomplete text`, async () => { + await assert.rejects( + read(id), + (error: unknown) => + error instanceof NativeSearchError && /live text extraction limits/.test(error.message) + ) + }) + } + await check( + 'DOCX callers without complete mode retain their existing extraction behavior', + async () => { + const expected = '🙂'.repeat(17) + const bytes = await docx(paragraph(expected)) + const result = await new DocxParser().parseBuffer(bytes, { + contentMode: 'complete', + maxTextBytes: 64, + }) + assert.equal(result.content, expected) + } + ) + await check('DOCX extraction honors a caller UTF-8 text budget', async () => { + const bytes = await docx(paragraph('🙂'.repeat(17))) + await assert.rejects( + new DocxParser().parseBuffer(bytes, { docxTextMode: 'complete', maxTextBytes: 64 }), + { + code: 'complexity_limit', + } + ) + }) + + for (const [id, metadata] of [ + ['download-denied', { capabilities: { canDownload: false } }], + ['download-unknown', { capabilities: {} }], + ['wrong-identity', { id: 'another-synthetic-file' }], + ['declared-oversized', { size: String(4 * MIB + 1) }], + ] as const) { + fixtures.set(id, { mimeType: PDF_MIME, body: pdfBody, metadata }) + await check(`${id} rejects before downloading file bytes`, async () => { + await unavailable(id) + assert.equal( + requests.filter((request) => request.id === id && request.resource === 'media').length, + 0 + ) + }) + } + + const negativeFixtures: [string, Fixture][] = [ + [ + 'pdf-text-mismatch', + { + mimeType: PDF_MIME, + body: Buffer.from('Provider diagnostic text must not become PDF content'), + }, + ], + [ + 'docx-html-mismatch', + { + mimeType: DOCX_MIME, + body: Buffer.from( + 'Provider error page must not become DOCX content' + ), + }, + ], + ['cross-format-mismatch', { mimeType: PDF_MIME, body: wordBody }], + ['malformed-pdf', { mimeType: PDF_MIME, body: Buffer.from('%PDF-1.4\nnot a valid document') }], + ['malformed-docx', { mimeType: DOCX_MIME, body: Buffer.from('PK\x03\x04broken archive') }], + ['encrypted-pdf', { mimeType: PDF_MIME, body: encryptedPdf() }], + ['text-free-pdf', { mimeType: PDF_MIME, body: await pdf(['']) }], + ['text-free-docx', { mimeType: DOCX_MIME, body: await docx('') }], + [ + 'too-many-pdf-pages', + { + mimeType: PDF_MIME, + body: await pdf( + Array.from({ length: 101 }, (_, index) => + index === 0 ? 'The synthetic page limit evidence remains readable.' : '' + ) + ), + }, + ], + [ + 'pdf-output-cap', + { + mimeType: PDF_MIME, + body: await pdf( + Array.from({ length: 90 }, (_, page) => + Array.from( + { length: 80 }, + (_, line) => `${page}-${line} ${'budget evidence '.repeat(12)}` + ).join('\n') + ) + ), + }, + ], + [ + 'docx-output-cap', + { mimeType: DOCX_MIME, body: await docx(paragraph(randomBytes(790_000).toString('base64'))) }, + ], + [ + 'docx-entry-cap', + { + mimeType: DOCX_MIME, + body: await docx(paragraph('Small synthetic body'), { + 'word/media/oversized.bin': Buffer.concat([ + randomBytes(2 * MIB), + Buffer.alloc(2 * MIB + 1), + ]), + }), + }, + ], + [ + 'docx-total-cap', + { + mimeType: DOCX_MIME, + body: await docx( + paragraph('Small synthetic body'), + Object.fromEntries( + [1, 2, 3].map((index) => [ + `word/media/part-${index}.bin`, + Buffer.concat([randomBytes(MIB), Buffer.alloc(4 * MIB - 1024 - MIB)]), + ]) + ) + ), + }, + ], + ] + for (const [id, fixture] of negativeFixtures) { + fixtures.set(id, fixture) + await check(`${id} rejects instead of returning metadata as extracted content`, () => + unavailable(id) + ) + } + + fixtures.set('stream-overflow', { + mimeType: PDF_MIME, + body: pdfBody, + mode: 'overflow', + metadata: { size: undefined }, + }) + await check('chunked media overflow reports an actionable unavailable result', async () => { + await assert.rejects( + read('stream-overflow'), + (error: unknown) => + error instanceof NativeSearchError && + error.status === 'unavailable' && + /4 MiB/.test(error.message) + ) + assert.equal( + requests.filter((request) => request.id === 'stream-overflow' && request.resource === 'media') + .length, + 1 + ) + }) + + fixtures.set('cancel-stream', { mimeType: PDF_MIME, body: pdfBody, mode: 'stall' }) + await check( + 'cancelling a media read closes its upstream and releases no partial document', + async () => { + const controller = new AbortController() + const reading = read( + 'cancel-stream', + AbortSignal.any([controller.signal, AbortSignal.timeout(5000)]) + ) + try { + await Promise.race([ + mediaStarted.promise, + reading.then(() => assert.fail('A PDF read returned without requesting its bytes')), + ]) + controller.abort(new DOMException('Synthetic read cancelled', 'AbortError')) + await assert.rejects( + reading, + (error: unknown) => error instanceof Error && /abort|cancel/i.test(error.message) + ) + await Promise.race([ + mediaClosed.promise, + new Promise((_resolve, reject) => + AbortSignal.timeout(1000).addEventListener( + 'abort', + () => reject(new Error('Cancelled media connection stayed open')), + { once: true } + ) + ), + ]) + } finally { + controller.abort() + } + } + ) + + for (const { name, pages, expected, budget } of [ + { + name: 'ASCII text below its byte budget', + pages: ['A'.repeat(60)], + expected: 'A'.repeat(60), + budget: 64, + }, + { + name: 'ASCII text exactly at its byte budget', + pages: ['B'.repeat(64)], + expected: 'B'.repeat(64), + budget: 64, + }, + { + name: 'UTF-8 text below its byte budget', + pages: ['€'.repeat(20)], + expected: '€'.repeat(20), + budget: 64, + }, + { + name: 'normalized text at its byte budget', + pages: [' Text '], + expected: 'Text', + budget: 4, + }, + { + name: 'page separators at the byte budget', + pages: ['A'.repeat(30), 'B'.repeat(32)], + expected: `${'A'.repeat(30)}\n\n${'B'.repeat(32)}`, + budget: 64, + }, + ]) { + await check(`PDF complete extraction accepts ${name}`, async () => { + const result = await new PdfParser().parseBuffer(await pdf(pages), { + pdfTextMode: 'complete', + maxTextBytes: budget, + }) + assert.equal(result.content, expected) + }) + } + await check('PDF output byte budget includes separators between rendered pages', async () => { + const pages = await pdf(['A'.repeat(30), 'B'.repeat(32)]) + await assert.rejects( + new PdfParser().parseBuffer(pages, { pdfTextMode: 'complete', maxTextBytes: 63 }), + { code: 'complexity_limit' } + ) + }) + + await check('PDF complete extraction enforces the caller page budget', async () => { + await assert.rejects( + new PdfParser().parseBuffer(pdfBody, { pdfTextMode: 'complete', pdfMaxPages: 1 }), + { code: 'complexity_limit' } + ) + }) + await check( + 'PDF complete extraction counts UTF-8 bytes against the caller output budget', + async () => { + const multibyte = await pdf(['€'.repeat(30)]) + await assert.rejects( + new PdfParser().parseBuffer(multibyte, { pdfTextMode: 'complete', maxTextBytes: 64 }), + { code: 'complexity_limit' } + ) + } + ) + await check('an already-cancelled read makes no provider request', async () => { + const before = requests.length + const controller = new AbortController() + controller.abort() + await assert.rejects(read('pdf', controller.signal)) + assert.equal(requests.length, before) + }) +} catch (error) { + checks.push({ + name: 'Synthetic HTTP acceptance setup', + status: 'failed', + durationMs: 0, + error: truncate(getErrorMessage(error), 1000), + }) +} finally { + server.closeAllConnections() + await new Promise((resolve) => server.close(() => resolve())) + await mkdir(dirname(reportPath), { recursive: true }) + await writeFile( + reportPath, + JSON.stringify( + { + checkedAt: new Date().toISOString(), + boundary: + 'Real loopback HTTP native client, Drive adapter and document parsers; synthetic fixtures; no application authorization or connected Google account.', + checks, + requests, + }, + null, + 2 + ), + { mode: 0o600 } + ) +} +const failed = checks.filter((check) => check.status === 'failed').length +logger.info('Google content acceptance completed', { + passed: checks.length - failed, + failed, + reportPath, +}) +if (failed) process.exitCode = 1