diff --git a/.dockerignore b/.dockerignore index 4e7bb6db..288e2245 100644 --- a/.dockerignore +++ b/.dockerignore @@ -2,4 +2,7 @@ __pycache__ uploads/ myenv/ -venv/ \ No newline at end of file +venv/ +.git +.worktrees +node_modules diff --git a/.github/workflows/bun-service.yml b/.github/workflows/bun-service.yml new file mode 100644 index 00000000..6c1d4871 --- /dev/null +++ b/.github/workflows/bun-service.yml @@ -0,0 +1,32 @@ +name: Bun service + +on: + pull_request: + push: + branches: [main] + +jobs: + extraction: + runs-on: ubuntu-latest + defaults: + run: + working-directory: service + steps: + - uses: actions/checkout@v4 + - uses: oven-sh/setup-bun@v2 + with: + bun-version: '1.4.2' + - name: Install locked dependencies + run: bun install --frozen-lockfile + - name: Typecheck + run: bun run typecheck + - name: Formatting + run: bun run format:check + - name: Native extraction and service contract tests + run: bun test + - name: Build standalone image + run: docker build -f service/Dockerfile -t rag-bun-test . + working-directory: . + - name: Verify image starts without Python, database or provider credentials + run: bash service/test/image.sh rag-bun-test + working-directory: . diff --git a/.gitignore b/.gitignore index 38921790..57141d6e 100644 --- a/.gitignore +++ b/.gitignore @@ -9,3 +9,5 @@ venv/ *.pyc dev.yml SHOPIFY.md + +node_modules/ diff --git a/EXTRACTION.md b/EXTRACTION.md new file mode 100644 index 00000000..5ab576bc --- /dev/null +++ b/EXTRACTION.md @@ -0,0 +1,141 @@ +# Bun/Hono RAG service: first extraction slice + +PR #330 starts the **new Bun/Hono service**, under `service/`. It is not an +extension of FastAPI and does not proxy to Python. The existing Python API, +requirements, ingestion, `/text`, collection schema and deployment remain +unchanged. Keep that service available throughout compatibility evaluation. + +This first slice implements only `POST /v1/extract` with the `document-v1` DOCX +profile. Retrieval, embeddings, reranking, PDF accounting, OCR and the LibreChat +client migration remain follow-up slices. **Do not point LibreChat's existing +`RAG_API_URL` at this service yet:** it does not implement the legacy endpoints. + +## Run independently + +Requires Bun 1.4.2. The Node AnyDoc binding is pinned to 0.1.3, matching Marco's +LibreChat PR #14701. The binding runs in a separate, killable Bun process, never +in the serving process. No Node server, Python runtime, database or embedding +provider is required. + +```sh +cd service +bun install --frozen-lockfile +bun run typecheck +bun test +bun run start +``` + +The default port is **8001**, so the existing Python service may stay on 8000. +`GET /health` works without external dependencies. Extraction is disabled by +default, and `/v1/extract` is not registered when off. + +Build a separate image from the repository root: + +```sh +docker build -f service/Dockerfile -t rag-bun . +``` + +Opt in at startup with `RAG_EXTRACTION_API_ENABLED=true` and `RAG_JWT_SECRET` +(minimum 32 characters). The secret must be the dedicated RAG signing key, **not +LibreChat's session key**. Each extraction request needs an HS256 JWT with +`sub`, `exp`, issuer `librechat`, audience `rag-api`, and a `scopes` array +containing `rag:documents`. Issuer and audience are configurable. Legacy `{id}` +tokens and inference-only tokens do not authorize this new route. That is a +new-service contract, not a change to the old Python API. Coordinate the +LibreChat token supplier with the strict-auth migration before sending traffic. + +## Request and response + +Send multipart `profile=document-v1` and one `file`. DOCX MIME is authoritative; +a generic MIME requires a `.docx` filename. No caller-provided path or URL is +read, and no hosted parser, OCR service, or inference provider is called. + +Successful response: + +```json +{ + "profile": "document-v1", + "text": "# Extracted Markdown\n", + "format": "markdown", + "completeness": "complete", + "may_omit_content": false, + "pages_needing_ocr": [], + "truncated": false, + "parser": { "name": "anydoc", "version": "0.1.3" } +} +``` + +An archive with non-thumbnail artwork or embedded objects is marked `partial` +and `may_omit_content=true`. Package relationships and content types identify +artwork even when its filename has no image extension. The root relationship +resolves the main document; a conventional `word/document.xml` path is not +required. This is conservative omission detection, **not a +proof that every source element is inspectable**. Never use partial text or a +preview as complete content-inspection input. No automatic fallback occurs +inside the service, so hard refusals cannot accidentally become paid OCR calls. + +`detail.code` classifies failures without source text, tokens, or native error +messages: + +| Status | Codes | Meaning | +|---|---|---| +| 400/415 | `INVALID_MULTIPART`, `UNSUPPORTED_PROFILE`, `UNSUPPORTED_DOCUMENT_TYPE` | Malformed or unsupported request | +| 401/403 | `EXTRACTION_AUTH_REQUIRED`, `EXTRACTION_FORBIDDEN` | Missing/invalid service token or missing document scope | +| 404 | `EXTRACTION_DISABLED` | Extraction not registered | +| 413 | `PARSER_INPUT_LIMIT`, `PARSER_OUTPUT_LIMIT`, `ZIP_BOMB` | Hard refusal; do not send those bytes to another parser | +| 422 | `ARCHIVE_INVALID`, `NO_DOCUMENT_TEXT`, `PARSE_FAILED` | Unusable archive or no usable conversion | +| 429 | `CONCURRENCY_LIMIT` | Busy; retry later rather than escalating to OCR | +| 503/504 | `PARSER_UNAVAILABLE`, `PARSER_CRASH`, `PARSER_TIMEOUT` | Infrastructure failure or overall deadline | +| 408 | `REQUEST_CANCELLED` | Cancelled operation; no result retained | + +## Limits and lifecycle + +Authentication and parser admission happen **before** multipart consumption. +Multipart uploads stream to a private random directory. Both declared length +and actual streamed bytes are bounded; there is no `request.formData()` buffer +or health-check round trip. The Bun listener also enforces the body ceiling. + +| Environment variable | Default | +|---|---:| +| `RAG_PORT` / `RAG_HOST` | 8001 / 0.0.0.0 | +| `RAG_EXTRACTION_CONCURRENT` | 2 | +| `RAG_EXTRACTION_QUEUED` | 6 | +| `RAG_EXTRACTION_TIMEOUT_MS` | 30,000 | +| `RAG_EXTRACTION_MAX_FILE_BYTES` | 15 MiB | +| `RAG_EXTRACTION_MAX_BODY_BYTES` | 16 MiB | +| `RAG_EXTRACTION_MAX_OUTPUT_BYTES` | 15 MiB | +| `RAG_EXTRACTION_MAX_ENTRY_BYTES` | 25 MiB | +| `RAG_EXTRACTION_MAX_ARCHIVE_BYTES` | 100 MiB | +| `RAG_EXTRACTION_MAX_ENTRIES` | 4,096 | +| `RAG_EXTRACTION_TEMP_DIR` | operating-system temp directory | +| `RAG_JWT_ISSUER` / `RAG_JWT_AUDIENCE` | librechat / rag-api | + +Limits are per service process, not cluster-wide quotas. Queue wait, upload +staging and parsing share one deadline. Cancelled queued work is removed; an +active child is killed and reaped before cleanup and slot reuse. The child +receives no JWT secret or provider credentials, and Bun's automatic dotenv +loading is disabled in it. Both sides cap serialized IPC +output. Actual decompressed entry bytes are checked before native parsing; +metadata-only size claims are not trusted. Graceful shutdown stops accepting +requests and allows active operations to finish within their deadlines. + +## Verification and adoption + +The Bun corpus uses Marco's structured DOCX at +`fb7bbcd9cf75f4f78ecbd5a8780685c481600be2` and the exact output expected by the +Python prototype. Tests use real native parsing, Hono requests and child +processes, with injected crash/hang/overproduction programs for failure cases. +Typecheck and native tests run in a separate Bun CI job alongside the unchanged +Python jobs. + +No traffic cutover or database migration is part of this PR. Next, add a flagged +LibreChat adapter using this result contract, then prove content inspection, +preview/sharing, token scopes, outages and mixed-version behavior end to end. +Keep raw Markdown, rich HTML preview and complete inspection as distinct +contracts. Expand formats or enable reranking only behind their own fidelity +and quality gates. No measured speedup or full service parity is claimed here. + +The listener idle timeout is set above the overall extraction deadline so a +parse lasting more than Bun's default ten seconds is not reset prematurely. +`RAG_EXTRACTION_TIMEOUT_MS` may be at most 254,000, leaving one second within +Bun's 255-second listener limit for the result or typed timeout response. diff --git a/service/Dockerfile b/service/Dockerfile new file mode 100644 index 00000000..27985f22 --- /dev/null +++ b/service/Dockerfile @@ -0,0 +1,13 @@ +FROM oven/bun:1.4.2-slim AS dependencies +WORKDIR /app +COPY service/package.json service/bun.lock ./ +RUN bun install --frozen-lockfile --production + +FROM oven/bun:1.4.2-slim +WORKDIR /app +COPY --from=dependencies --chown=bun:bun /app/node_modules ./node_modules +COPY --chown=bun:bun service/package.json ./package.json +COPY --chown=bun:bun service/src ./src +USER bun +EXPOSE 8001 +CMD ["bun", "run", "src/server.ts"] diff --git a/service/bun.lock b/service/bun.lock new file mode 100644 index 00000000..d10a5005 --- /dev/null +++ b/service/bun.lock @@ -0,0 +1,79 @@ +{ + "lockfileVersion": 2, + "configVersion": 1, + "workspaces": { + "": { + "name": "@librechat/rag-service", + "dependencies": { + "@firecrawl/anydoc": "0.1.3", + "busboy": "1.6.0", + "hono": "4.13.11", + "jose": "6.2.12", + "saxes": "6.0.0", + "yauzl": "3.4.0", + "zod": "4.3.6", + }, + "devDependencies": { + "@types/bun": "1.4.2", + "@types/busboy": "1.5.4", + "@types/yauzl": "3.4.0", + "fflate": "0.8.2", + "prettier": "3.8.1", + "typescript": "5.9.3", + }, + }, + }, + "packages": { + "@firecrawl/anydoc": ["@firecrawl/anydoc@0.1.3", "", { "optionalDependencies": { "@firecrawl/anydoc-darwin-arm64": "0.1.3", "@firecrawl/anydoc-darwin-x64": "0.1.3", "@firecrawl/anydoc-linux-arm64-gnu": "0.1.3", "@firecrawl/anydoc-linux-arm64-musl": "0.1.3", "@firecrawl/anydoc-linux-x64-gnu": "0.1.3", "@firecrawl/anydoc-linux-x64-musl": "0.1.3", "@firecrawl/anydoc-win32-x64-msvc": "0.1.3" }, "bin": { "anydoc": "cli.js" } }, "sha512-OuYZWJHxiSNxU414g3l4w0gKb5BTqVZKGCwgmrkn8Mq84lf5z1Y2pHcTWQOWbyN4qD2NX7hQhrnozGhYYql6sw=="], + + "@firecrawl/anydoc-darwin-arm64": ["@firecrawl/anydoc-darwin-arm64@0.1.3", "", { "os": "darwin", "cpu": "arm64" }, "sha512-tCi0pqATVbS4el+lF/XBGvXzMUzc1dLTQd6eZxctLoP5HBLyQsmBQiSDdU+dE5reVpriLdumAo3MRhSJVBfikw=="], + + "@firecrawl/anydoc-darwin-x64": ["@firecrawl/anydoc-darwin-x64@0.1.3", "", { "os": "darwin", "cpu": "x64" }, "sha512-WpegrFgCZfXgbDBsrHrZcUcJNoB31fL7ztGL0x3OCThGNcND/JFOAb0ihS3Y/uwbHvXJDZTt+ZzAQOhSoC/8Rg=="], + + "@firecrawl/anydoc-linux-arm64-gnu": ["@firecrawl/anydoc-linux-arm64-gnu@0.1.3", "", { "os": "linux", "cpu": "arm64" }, "sha512-5/QCKECO2ujVZ7b3JOcNyvZrH1uSB3I/uLmgN05FWEIERA3rrwdzZi/k6K7ULcmSOD4eMJDfd6g41Ka48gW6zw=="], + + "@firecrawl/anydoc-linux-arm64-musl": ["@firecrawl/anydoc-linux-arm64-musl@0.1.3", "", { "os": "linux", "cpu": "arm64" }, "sha512-nf/IakkpFAA3wdIoicohgOW+3D3xDr8gGhU5C+QdWDvRTQjK3Bm2R0DLvZJcnuOonsb01ob+BqFb9Q7PqmCOpw=="], + + "@firecrawl/anydoc-linux-x64-gnu": ["@firecrawl/anydoc-linux-x64-gnu@0.1.3", "", { "os": "linux", "cpu": "x64" }, "sha512-VYnBbTIyRwv3HcWZgEaguPMjRs2CdA5aV+Nx60kCBV4Ykp3GOIdiCxJFE0RfOUQN3nfPV6hcaKtXz+1O4jxX0A=="], + + "@firecrawl/anydoc-linux-x64-musl": ["@firecrawl/anydoc-linux-x64-musl@0.1.3", "", { "os": "linux", "cpu": "x64" }, "sha512-gL3dYfVmpN6ixfdmy2Ifgkdd5s3xrIHuwUlTTsTWR15OkhfL9LuA/VkgDO+p82D3SWAVnCTEWI7scE8nkMLmBQ=="], + + "@firecrawl/anydoc-win32-x64-msvc": ["@firecrawl/anydoc-win32-x64-msvc@0.1.3", "", { "os": "win32", "cpu": "x64" }, "sha512-NwjmXqhrL6xVTLAJeZ0Ie7h7Clm2CwSPuJ9wuZX9fPP5nx/p7GvPRDvtLq6PxtXk2LQFWklfHmDczUQmzMh1Kg=="], + + "@types/bun": ["@types/bun@1.4.2", "", { "dependencies": { "bun-types": "1.4.2" } }, "sha512-GimotNn7+ZV0uVArItBbriZsR1oNf0+WTzPkdcFrzShI7k2norL0uzEaJT8T33dWr7O/c9ZDuAFQrctKCi72oQ=="], + + "@types/busboy": ["@types/busboy@1.5.4", "", { "dependencies": { "@types/node": "*" } }, "sha512-kG7WrUuAKK0NoyxfQHsVE6j1m01s6kMma64E+OZenQABMQyTJop1DumUWcLwAQ2JzpefU7PDYoRDKl8uZosFjw=="], + + "@types/node": ["@types/node@26.6.3", "", { "dependencies": { "undici-types": "~8.9.0" } }, "sha512-dsqMQQoeTLqu9wynDD00q573mNzso3IdQOAfHRJqLCcmCFPoGo9A1bDpUcv/9tnKpErQWv9uKeGfl37EIS02Yg=="], + + "@types/yauzl": ["@types/yauzl@3.4.0", "", { "dependencies": { "@types/node": "*" } }, "sha512-NRPn5w6h8dhcnmx3YIRQcqMywY/+nND/uOkJessedcrowO3C0AssHp3tMJpxKAwOhFOo0OV1y9VtsC5hbKKBAw=="], + + "bun-types": ["bun-types@1.4.2", "", { "dependencies": { "@types/node": "*" } }, "sha512-bxV1FgK7yBIzjRe5zBozIM4Bem11ZJcCXSrjWRG3YWLt8yFDePu4cLjpebO8OvPeIE9trbyPF4fuj3Cia4Fj3w=="], + + "busboy": ["busboy@1.6.0", "", { "dependencies": { "streamsearch": "^1.1.0" } }, "sha512-8SFQbg/0hQ9xy3UNTB0YEnsNBbWfhf7RtnzpL7TkBiTBRfrQ9Fxcnz7VJsleJpyp6rVLvXiuORqjlHi5q+PYuA=="], + + "fflate": ["fflate@0.8.2", "", {}, "sha512-cPJU47OaAoCbg0pBvzsgpTPhmhqI5eJjh/JIu8tPj5q+T7iLvW/JAYUqmE7KOB4R1ZyEhzBaIQpQpardBF5z8A=="], + + "hono": ["hono@4.13.11", "", {}, "sha512-/SMX/RQNJn7oNmFwH6DtwDqcrzZU2otm1FD6OJe2cWF6cfy/F8hq0XVBv7xW3FTJshRQwuBnXHDq2kNsdZnshg=="], + + "jose": ["jose@6.2.12", "", {}, "sha512-9NiFmJEex0sy2Dk58j2UGBSHgUs2ypF9eZSu4L6vjOX3Dp96Sw1F3uL+H+D1sx02jZZdzUT0HgvCy59CuvXcWw=="], + + "pend": ["pend@1.2.0", "", {}, "sha512-F3asv42UuXchdzt+xXqfW1OGlVBe+mxa2mqI0pg5yAHZPvFmY3Y6drSf/GQ1A86WgWEN9Kzh/WrgKa6iGcHXLg=="], + + "prettier": ["prettier@3.8.1", "", { "bin": { "prettier": "bin/prettier.cjs" } }, "sha512-UOnG6LftzbdaHZcKoPFtOcCKztrQ57WkHDeRD9t/PTQtmT0NHSeWWepj6pS0z/N7+08BHFDQVUrfmfMRcZwbMg=="], + + "saxes": ["saxes@6.0.0", "", { "dependencies": { "xmlchars": "^2.2.0" } }, "sha512-xAg7SOnEhrm5zI3puOOKyy1OMcMlIJZYNJY7xLBwSze0UjhPLnWfj2GF2EpT0jmzaJKIWKHLsaSSajf35bcYnA=="], + + "streamsearch": ["streamsearch@1.1.0", "", {}, "sha512-Mcc5wHehp9aXz1ax6bZUyY5afg9u2rv5cqQI3mRrYkGC8rW2hM02jWuwjtL++LS5qinSyhj2QfLyNsuc+VsExg=="], + + "typescript": ["typescript@5.9.3", "", { "bin": { "tsc": "bin/tsc", "tsserver": "bin/tsserver" } }, "sha512-jl1vZzPDinLr9eUt3J/t7V6FgNEw9QjvBPdysz9KfQDD41fQrC2Y4vKQdiaUpFT4bXlb1RHhLpp8wtm6M5TgSw=="], + + "undici-types": ["undici-types@8.9.0", "", {}, "sha512-KTDyRTYX8sWmKXAikPHHSyc63CRPETMctyjKFupcC6OBLXT3xsN0e9aF7m+mIXutFWpUXuedtowG7iLOzp0kQg=="], + + "xmlchars": ["xmlchars@2.2.0", "", {}, "sha512-JZnDKK8B0RCDw84FNdDAIpZK+JuJw+s7Lz8nksI7SIuU3UXJJslUthsi+uWBUYOwPFwW7W7PRLRfUKpxjtjFCw=="], + + "yauzl": ["yauzl@3.4.0", "", { "dependencies": { "pend": "~1.2.0" } }, "sha512-jIH9yLR9wqr0wOS0TpBvo/g/2UgZH5qePVbjgRliiF0BYvOZyaBknKsF+x9Iht0O6sqgnB93rCICdOZFecJuDw=="], + + "zod": ["zod@4.3.6", "", {}, "sha512-rftlrkhHZOcjDwkGlnUtZZkvaPHCsDATp4pGpuOOMDaTdDDXF91wuVDJoWoPsKX/3YPQ5fHuF3STjcYyKr+Qhg=="], + } +} diff --git a/service/package.json b/service/package.json new file mode 100644 index 00000000..b38f9ca5 --- /dev/null +++ b/service/package.json @@ -0,0 +1,31 @@ +{ + "name": "@librechat/rag-service", + "private": true, + "version": "0.1.0", + "type": "module", + "packageManager": "bun@1.4.2", + "scripts": { + "start": "bun run src/server.ts", + "test": "bun test", + "typecheck": "tsc --noEmit", + "format:check": "prettier --check src test *.json", + "format": "prettier --write src test *.json" + }, + "dependencies": { + "@firecrawl/anydoc": "0.1.3", + "busboy": "1.6.0", + "hono": "4.13.11", + "jose": "6.2.12", + "saxes": "6.0.0", + "yauzl": "3.4.0", + "zod": "4.3.6" + }, + "devDependencies": { + "@types/bun": "1.4.2", + "@types/busboy": "1.5.4", + "@types/yauzl": "3.4.0", + "fflate": "0.8.2", + "prettier": "3.8.1", + "typescript": "5.9.3" + } +} diff --git a/service/src/admission.ts b/service/src/admission.ts new file mode 100644 index 00000000..ece93250 --- /dev/null +++ b/service/src/admission.ts @@ -0,0 +1,56 @@ +import { ExtractionError } from "./contract"; + +type Release = () => void; +type Waiter = { + signal: AbortSignal; + resolve: (release: Release) => void; + reject: (error: Error) => void; + abort: () => void; +}; + +export class Admission { + private active = 0; + private readonly waiting: Waiter[] = []; + constructor( + private readonly concurrent: number, + private readonly queued: number, + ) {} + + acquire(signal: AbortSignal): Promise { + signal.throwIfAborted(); + if (this.active < this.concurrent) { + this.active++; + return Promise.resolve(this.release()); + } + if (this.waiting.length >= this.queued) { + return Promise.reject(new ExtractionError("CONCURRENCY_LIMIT")); + } + return new Promise((resolve, reject) => { + const waiter: Waiter = { + signal, + resolve, + reject, + abort: () => { + const index = this.waiting.indexOf(waiter); + if (index >= 0) this.waiting.splice(index, 1); + reject(new ExtractionError("REQUEST_CANCELLED")); + }, + }; + signal.addEventListener("abort", waiter.abort, { once: true }); + this.waiting.push(waiter); + }); + } + + private release(): Release { + let released = false; + return () => { + if (released) return; + released = true; + const next = this.waiting.shift(); + if (next) { + next.signal.removeEventListener("abort", next.abort); + next.resolve(this.release()); + } else this.active--; + }; + } +} diff --git a/service/src/app.ts b/service/src/app.ts new file mode 100644 index 00000000..72453bd5 --- /dev/null +++ b/service/src/app.ts @@ -0,0 +1,128 @@ +import { Hono } from "hono"; +import { jwtVerify } from "jose"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import type { ContentfulStatusCode } from "hono/utils/http-status"; +import type { Runner } from "./process"; +import { Admission } from "./admission"; +import { configSchema, type Config } from "./config"; +import { ExtractionError, type ErrorCode } from "./contract"; +import { runWorker } from "./process"; +import { stageUpload } from "./upload"; + +const status: Record = { + EXTRACTION_DISABLED: 404, + EXTRACTION_AUTH_REQUIRED: 401, + EXTRACTION_FORBIDDEN: 403, + UNSUPPORTED_PROFILE: 400, + UNSUPPORTED_DOCUMENT_TYPE: 415, + INVALID_MULTIPART: 400, + PARSER_INPUT_LIMIT: 413, + PARSER_OUTPUT_LIMIT: 413, + ZIP_BOMB: 413, + ARCHIVE_INVALID: 422, + NO_DOCUMENT_TEXT: 422, + PARSE_FAILED: 422, + CONCURRENCY_LIMIT: 429, + PARSER_CRASH: 503, + PARSER_UNAVAILABLE: 503, + PARSER_TIMEOUT: 504, + REQUEST_CANCELLED: 408, +}; + +export function createApp( + input: Partial = {}, + runner: Runner = runWorker, +): Hono { + const config = configSchema.parse(input); + const app = new Hono(); + const admission = new Admission(config.concurrent, config.queued); + app.onError((error, context) => { + const code = + error instanceof ExtractionError ? error.code : "PARSER_UNAVAILABLE"; + return context.json({ detail: { code } }, status[code]); + }); + app.notFound((context) => + context.json({ detail: { code: "EXTRACTION_DISABLED" } }, 404), + ); + app.get("/health", (context) => + context.json({ + status: "UP", + service: "rag-bun", + extraction_profiles: config.enabled ? ["document-v1"] : [], + }), + ); + if (!config.enabled) return app; + const key = new TextEncoder().encode(config.secret); + app.post("/v1/extract", async (context) => { + const authorization = context.req.header("Authorization"); + if (!authorization?.startsWith("Bearer ")) + throw new ExtractionError("EXTRACTION_AUTH_REQUIRED"); + try { + const { payload } = await jwtVerify(authorization.slice(7), key, { + algorithms: ["HS256"], + issuer: config.issuer, + audience: config.audience, + requiredClaims: ["exp", "sub"], + }); + if (typeof payload.sub !== "string" || !payload.sub) + throw new ExtractionError("EXTRACTION_AUTH_REQUIRED"); + if ( + !Array.isArray(payload.scopes) || + !payload.scopes.includes("rag:documents") + ) { + throw new ExtractionError("EXTRACTION_FORBIDDEN"); + } + } catch (error) { + if (error instanceof ExtractionError) throw error; + throw new ExtractionError("EXTRACTION_AUTH_REQUIRED"); + } + const length = context.req.header("Content-Length"); + if ( + length !== undefined && + (!/^\d+$/.test(length) || Number(length) > config.maxBodyBytes) + ) { + throw new ExtractionError("PARSER_INPUT_LIMIT"); + } + const deadline = new AbortController(); + const timer = setTimeout(() => deadline.abort(), config.timeoutMs); + const signal = AbortSignal.any([context.req.raw.signal, deadline.signal]); + let release: (() => void) | undefined; + let directory: string | undefined; + try { + release = await admission.acquire(signal); + signal.throwIfAborted(); + directory = await mkdtemp( + join(config.tempRoot ?? tmpdir(), "rag-extract-"), + ); + signal.throwIfAborted(); + const path = join(directory, "input.docx"); + await stageUpload(context.req.raw, path, config, signal); + const result = await runner( + { + path, + maxOutputBytes: config.maxOutputBytes, + maxEntryBytes: config.maxEntryBytes, + maxArchiveBytes: config.maxArchiveBytes, + maxEntries: config.maxEntries, + }, + signal, + ); + signal.throwIfAborted(); + return context.json(result); + } catch (error) { + if (deadline.signal.aborted) throw new ExtractionError("PARSER_TIMEOUT"); + if (signal.aborted) throw new ExtractionError("REQUEST_CANCELLED"); + throw error; + } finally { + clearTimeout(timer); + try { + if (directory) await rm(directory, { recursive: true, force: true }); + } finally { + release?.(); + } + } + }); + return app; +} diff --git a/service/src/config.ts b/service/src/config.ts new file mode 100644 index 00000000..512fc54c --- /dev/null +++ b/service/src/config.ts @@ -0,0 +1,67 @@ +import { z } from "zod"; + +const positive = z.number().int().positive(); +export const configSchema = z + .object({ + enabled: z.boolean().default(false), + secret: z.string().min(32).optional(), + issuer: z.string().min(1).default("librechat"), + audience: z.string().min(1).default("rag-api"), + concurrent: positive.default(2), + queued: z.number().int().nonnegative().default(6), + timeoutMs: positive.max(254_000).default(30_000), + maxFileBytes: positive.default(15 * 1024 * 1024), + maxBodyBytes: positive.default(16 * 1024 * 1024), + maxOutputBytes: positive.default(15 * 1024 * 1024), + maxEntryBytes: positive.default(25 * 1024 * 1024), + maxArchiveBytes: positive.default(100 * 1024 * 1024), + maxEntries: positive.default(4096), + tempRoot: z.string().min(1).optional(), + }) + .superRefine((config, ctx) => { + if (config.enabled && !config.secret) { + ctx.addIssue({ + code: "custom", + message: "Enabled extraction requires RAG_JWT_SECRET (32+ characters)", + }); + } + if (config.maxBodyBytes <= config.maxFileBytes) { + ctx.addIssue({ + code: "custom", + message: + "Body ceiling must exceed the file ceiling for multipart framing", + }); + } + }); +export type Config = z.infer; + +export function fromEnv(env: NodeJS.ProcessEnv): Config { + if (env.RAG_JWT_SECRET && env.RAG_JWT_SECRET === env.JWT_SECRET) { + throw new Error("RAG_JWT_SECRET must differ from JWT_SECRET"); + } + const number = (name: string) => + env[name] === undefined ? undefined : Number(env[name]); + const enabled = env.RAG_EXTRACTION_API_ENABLED; + if ( + enabled !== undefined && + !["true", "false", "1", "0"].includes(enabled.toLowerCase()) + ) { + throw new Error("RAG_EXTRACTION_API_ENABLED must be true, false, 1 or 0"); + } + return configSchema.parse({ + enabled: enabled === "1" || enabled?.toLowerCase() === "true", + secret: env.RAG_JWT_SECRET, + issuer: env.RAG_JWT_ISSUER, + audience: env.RAG_JWT_AUDIENCE, + concurrent: number("RAG_EXTRACTION_CONCURRENT"), + queued: number("RAG_EXTRACTION_QUEUED"), + timeoutMs: number("RAG_EXTRACTION_TIMEOUT_MS"), + maxFileBytes: number("RAG_EXTRACTION_MAX_FILE_BYTES"), + maxBodyBytes: number("RAG_EXTRACTION_MAX_BODY_BYTES"), + maxOutputBytes: number("RAG_EXTRACTION_MAX_OUTPUT_BYTES"), + maxEntryBytes: number("RAG_EXTRACTION_MAX_ENTRY_BYTES"), + maxArchiveBytes: number("RAG_EXTRACTION_MAX_ARCHIVE_BYTES"), + maxEntries: number("RAG_EXTRACTION_MAX_ENTRIES"), + tempRoot: env.RAG_EXTRACTION_TEMP_DIR, + }); +} diff --git a/service/src/contract.ts b/service/src/contract.ts new file mode 100644 index 00000000..34e2142a --- /dev/null +++ b/service/src/contract.ts @@ -0,0 +1,55 @@ +import { z } from "zod"; + +export const resultSchema = z + .object({ + profile: z.literal("document-v1"), + text: z.string().min(1), + format: z.literal("markdown"), + completeness: z.enum(["complete", "partial"]), + may_omit_content: z.boolean(), + pages_needing_ocr: z.array(z.number().int().positive()), + truncated: z.literal(false), + parser: z.object({ + name: z.literal("anydoc"), + version: z.literal("0.1.3"), + }), + }) + .strict() + .superRefine((result, ctx) => { + if ((result.completeness === "partial") !== result.may_omit_content) { + ctx.addIssue({ code: "custom", message: "Inconsistent completeness" }); + } + }); +export type ExtractionResult = z.infer; + +export const errorSchema = z.enum([ + "EXTRACTION_DISABLED", + "EXTRACTION_AUTH_REQUIRED", + "EXTRACTION_FORBIDDEN", + "UNSUPPORTED_PROFILE", + "UNSUPPORTED_DOCUMENT_TYPE", + "INVALID_MULTIPART", + "PARSER_INPUT_LIMIT", + "PARSER_OUTPUT_LIMIT", + "ZIP_BOMB", + "ARCHIVE_INVALID", + "NO_DOCUMENT_TEXT", + "PARSE_FAILED", + "CONCURRENCY_LIMIT", + "PARSER_CRASH", + "PARSER_UNAVAILABLE", + "PARSER_TIMEOUT", + "REQUEST_CANCELLED", +]); +export type ErrorCode = z.infer; +export class ExtractionError extends Error { + constructor(readonly code: ErrorCode) { + super(code); + } +} +export const workerResponseSchema = z.discriminatedUnion("ok", [ + z.object({ ok: z.literal(true), result: resultSchema }), + z.object({ ok: z.literal(false), code: errorSchema }), +]); +export const DOCX_TYPE = + "application/vnd.openxmlformats-officedocument.wordprocessingml.document"; diff --git a/service/src/process.ts b/service/src/process.ts new file mode 100644 index 00000000..a606a9be --- /dev/null +++ b/service/src/process.ts @@ -0,0 +1,69 @@ +import { fileURLToPath } from "node:url"; +import type { WorkerRequest } from "./worker"; +import { + ExtractionError, + workerResponseSchema, + type ExtractionResult, +} from "./contract"; + +const workerPath = fileURLToPath(new URL("./worker.ts", import.meta.url)); +export type Runner = ( + request: WorkerRequest, + signal: AbortSignal, +) => Promise; + +export async function runWorker( + request: WorkerRequest, + signal: AbortSignal, + command: readonly string[] = [process.execPath, workerPath], +): Promise { + signal.throwIfAborted(); + const executable = command[0]; + if (!executable) throw new ExtractionError("PARSER_UNAVAILABLE"); + const child = Bun.spawn([executable, "--no-env-file", ...command.slice(1)], { + stdin: "pipe", + stdout: "pipe", + stderr: "ignore", + // Native parsing gets no signing key, provider credential or service config. + env: { PATH: process.env.PATH ?? "" }, + }); + const abort = () => { + child.kill("SIGKILL"); + }; + signal.addEventListener("abort", abort, { once: true }); + const reader = child.stdout.getReader(); + const chunks: Uint8Array[] = []; + let size = 0; + try { + signal.throwIfAborted(); + child.stdin.write(JSON.stringify(request)); + child.stdin.end(); + while (true) { + const part = await reader.read(); + if (part.done) break; + size += part.value.byteLength; + if (size > request.maxOutputBytes) + throw new ExtractionError("PARSER_OUTPUT_LIMIT"); + chunks.push(part.value); + } + const exit = await child.exited; + signal.throwIfAborted(); + if (exit !== 0) throw new ExtractionError("PARSER_CRASH"); + const result = workerResponseSchema.safeParse( + JSON.parse(Buffer.concat(chunks, size).toString("utf8")), + ); + if (!result.success) throw new ExtractionError("PARSER_CRASH"); + if (!result.data.ok) throw new ExtractionError(result.data.code); + return result.data.result; + } catch (error) { + child.kill("SIGKILL"); + await child.exited; + await reader.cancel().catch(() => {}); + if (signal.aborted) throw new ExtractionError("REQUEST_CANCELLED"); + if (error instanceof ExtractionError) throw error; + throw new ExtractionError("PARSER_CRASH"); + } finally { + signal.removeEventListener("abort", abort); + reader.releaseLock(); + } +} diff --git a/service/src/server.ts b/service/src/server.ts new file mode 100644 index 00000000..fbe136a3 --- /dev/null +++ b/service/src/server.ts @@ -0,0 +1,47 @@ +import type { Hono } from "hono"; +import type { Config } from "./config"; +import { createApp } from "./app"; +import { fromEnv } from "./config"; + +export function listen( + app: Hono, + config: Config, + port: number, + hostname: string, +) { + return Bun.serve({ + hostname, + port, + // Bun includes pending handlers in its idle timer (default: 10 seconds). + // Its maximum is 255 seconds; config bounds the request deadline to 254. + idleTimeout: Math.ceil(config.timeoutMs / 1000) + 1, + maxRequestBodySize: config.maxBodyBytes, + fetch: app.fetch, + }); +} + +if (import.meta.main) { + try { + const config = fromEnv(process.env); + const port = Number(process.env.RAG_PORT ?? 8001); + if (!Number.isInteger(port) || port < 1 || port > 65535) + throw new Error("Invalid port"); + const server = listen( + createApp(config), + config, + port, + process.env.RAG_HOST ?? "0.0.0.0", + ); + for (const event of ["SIGINT", "SIGTERM"] as const) { + process.once(event, () => { + void server.stop(false); + }); + } + console.info(`RAG Bun service listening on port ${server.port}`); + } catch { + console.error( + "Unable to start RAG Bun service: check runtime configuration", + ); + process.exitCode = 1; + } +} diff --git a/service/src/upload.ts b/service/src/upload.ts new file mode 100644 index 00000000..c425f2b6 --- /dev/null +++ b/service/src/upload.ts @@ -0,0 +1,128 @@ +import busboy from "busboy"; +import { createWriteStream } from "node:fs"; +import { basename, extname } from "node:path"; +import { Readable } from "node:stream"; +import { pipeline } from "node:stream/promises"; +import type { Config } from "./config"; +import { DOCX_TYPE, ExtractionError } from "./contract"; + +export async function stageUpload( + request: Request, + path: string, + config: Config, + signal: AbortSignal, +): Promise { + if (!request.body) throw new ExtractionError("INVALID_MULTIPART"); + let parser: ReturnType; + try { + parser = busboy({ + headers: { "content-type": request.headers.get("content-type") ?? "" }, + limits: { + fileSize: config.maxFileBytes, + files: 1, + fields: 1, + parts: 3, + fieldSize: 64, + }, + }); + } catch { + throw new ExtractionError("INVALID_MULTIPART"); + } + const reader = request.body.getReader(); + const abort = () => { + void reader.cancel().catch(() => {}); + }; + signal.addEventListener("abort", abort, { once: true }); + let bytes = 0; + const input = Readable.from( + (async function* () { + try { + while (true) { + signal.throwIfAborted(); + const part = await reader.read(); + if (part.done) break; + bytes += part.value.byteLength; + if (bytes > config.maxBodyBytes) + throw new ExtractionError("PARSER_INPUT_LIMIT"); + yield part.value; + } + } finally { + // Stop unread HTTP input before releasing the reader's lock. + await reader.cancel().catch(() => {}); + reader.releaseLock(); + } + })(), + ); + let profile: string | undefined; + let fileCount = 0; + let refusal: ExtractionError | undefined; + const writes: Promise[] = []; + const fail = (error: ExtractionError) => { + refusal ??= error; + parser.destroy(error); + }; + parser.on("field", (name, value, info) => { + if (name !== "profile" || info.valueTruncated || profile !== undefined) { + fail(new ExtractionError("INVALID_MULTIPART")); + } else profile = value; + }); + parser.on("file", (name, stream, info) => { + // Busboy destroys its current file stream when the parser is refused, even + // before a disk pipeline exists. Refused streams need an error listener too. + stream.on("error", () => {}); + fileCount++; + const extension = extname(basename(info.filename)).toLowerCase(); + const generic = [ + "application/octet-stream", + "binary/octet-stream", + ].includes(info.mimeType); + if (name !== "file" || fileCount !== 1) { + stream.resume(); + fail(new ExtractionError("INVALID_MULTIPART")); + return; + } + if (info.mimeType !== DOCX_TYPE && !(generic && extension === ".docx")) { + stream.resume(); + fail(new ExtractionError("UNSUPPORTED_DOCUMENT_TYPE")); + return; + } + stream.once("limit", () => fail(new ExtractionError("PARSER_INPUT_LIMIT"))); + const write = pipeline( + stream, + createWriteStream(path, { flags: "wx", mode: 0o600 }), + { signal }, + ); + // Attach rejection handling at creation, not only after the multipart parser + // finishes, so a storage failure terminates upload and never becomes unhandled. + writes.push( + write.catch((error: Error) => { + parser.destroy(error); + throw error; + }), + ); + void writes.at(-1)?.catch(() => {}); + }); + for (const event of ["filesLimit", "fieldsLimit", "partsLimit"] as const) { + parser.on(event, () => fail(new ExtractionError("INVALID_MULTIPART"))); + } + try { + await pipeline(input, parser, { signal }); + await Promise.all(writes); + signal.throwIfAborted(); + if (fileCount !== 1 || profile === undefined) + throw new ExtractionError("INVALID_MULTIPART"); + if (profile !== "document-v1") + throw new ExtractionError("UNSUPPORTED_PROFILE"); + } catch (error) { + input.destroy(); + parser.destroy(); + await reader.cancel().catch(() => {}); + await Promise.allSettled(writes); + if (refusal) throw refusal; + if (error instanceof ExtractionError) throw error; + if (signal.aborted) throw new ExtractionError("REQUEST_CANCELLED"); + throw new ExtractionError("INVALID_MULTIPART"); + } finally { + signal.removeEventListener("abort", abort); + } +} diff --git a/service/src/worker.ts b/service/src/worker.ts new file mode 100644 index 00000000..cee316db --- /dev/null +++ b/service/src/worker.ts @@ -0,0 +1,246 @@ +import { readFile } from "node:fs/promises"; +import { createRequire } from "node:module"; +import { posix } from "node:path"; +import type { Readable } from "node:stream"; +import { SaxesParser } from "saxes"; +import { open, type Entry } from "yauzl"; +import { z } from "zod"; +import type { ExtractionResult } from "./contract"; +import { ExtractionError } from "./contract"; + +const requestSchema = z.object({ + path: z.string(), + maxOutputBytes: z.number().int().positive(), + maxEntryBytes: z.number().int().positive(), + maxArchiveBytes: z.number().int().positive(), + maxEntries: z.number().int().positive(), +}); +export type WorkerRequest = z.infer; +const IMAGE = + /\.(?:jpe?g|png|gif|tiff?|bmp|webp|jp2|jpx|avif|heic|heif|emf|wmf|svg)$/i; +const PREVIEW = /^(?:docProps|Thumbnails)\//i; + +function mainPart(target: string): string { + // OPC targets are package URIs, never filesystem or network locations. + const decoded = decodeURIComponent(target); + if ( + /^[a-z]+:/i.test(decoded) || + decoded.includes("\\") || + decoded.split("/").includes("..") + ) { + throw new ExtractionError("ARCHIVE_INVALID"); + } + return posix.normalize("/" + decoded).slice(1); +} + +function inspectArchive(path: string, limits: WorkerRequest): Promise { + return new Promise((resolve, reject) => { + open( + path, + { lazyEntries: true, validateEntrySizes: true }, + (error, zip) => { + if (error || !zip) { + reject(new ExtractionError("ARCHIVE_INVALID")); + return; + } + let settled = false; + let active: Readable | undefined; + let total = 0; + let media = false; + let main: string | undefined; + const names = new Set(); + const imageParts = new Set(); + const imageExtensions = new Set(); + const fail = (code: "ARCHIVE_INVALID" | "ZIP_BOMB") => { + if (settled) return; + settled = true; + active?.destroy(); + zip.close(); + reject(new ExtractionError(code)); + }; + zip.on("error", () => fail("ARCHIVE_INVALID")); + zip.on("end", () => { + if (settled) return; + if (!names.has("[Content_Types].xml") || !main || !names.has(main)) { + fail("ARCHIVE_INVALID"); + return; + } + for (const name of names) { + if (PREVIEW.test(name)) continue; + if ( + imageParts.has(name) || + imageExtensions.has(posix.extname(name).slice(1).toLowerCase()) + ) + media = true; + } + settled = true; + resolve(media); + }); + if (zip.entryCount > limits.maxEntries) { + fail("ZIP_BOMB"); + return; + } + zip.on("entry", (entry: Entry) => { + if (settled) return; + if (names.has(entry.fileName)) { + fail("ARCHIVE_INVALID"); + return; + } + names.add(entry.fileName); + if (/\/$/.test(entry.fileName)) { + zip.readEntry(); + return; + } + if ( + entry.uncompressedSize > limits.maxEntryBytes || + total + entry.uncompressedSize > limits.maxArchiveBytes + ) { + fail("ZIP_BOMB"); + return; + } + if ( + (IMAGE.test(entry.fileName) && !PREVIEW.test(entry.fileName)) || + /\/embeddings\//i.test(entry.fileName) + ) + media = true; + const isTypes = entry.fileName === "[Content_Types].xml"; + const isRelationships = entry.fileName.endsWith(".rels"); + let metadata: SaxesParser<{ xmlns: true }> | undefined; + if (isTypes || isRelationships) { + metadata = new SaxesParser({ xmlns: true }); + metadata.on("error", () => fail("ARCHIVE_INVALID")); + // No DTD/entity expansion or external XML resources in the admission guard. + metadata.on("doctype", () => fail("ARCHIVE_INVALID")); + metadata.on("opentag", (tag) => { + const attr = (name: string) => tag.attributes[name]?.value; + if ( + isTypes && + attr("ContentType")?.toLowerCase().startsWith("image/") + ) { + const part = attr("PartName"); + const extension = attr("Extension"); + if (part) imageParts.add(mainPart(part)); + if (extension) imageExtensions.add(extension.toLowerCase()); + } + if (!isRelationships || tag.local !== "Relationship") return; + const type = attr("Type") ?? ""; + if (/\/(image|oleObject|package)$/.test(type)) media = true; + if ( + entry.fileName !== "_rels/.rels" || + !type.endsWith("/officeDocument") + ) + return; + if ( + main || + attr("TargetMode") === "External" || + !attr("Target") + ) { + fail("ARCHIVE_INVALID"); + return; + } + main = mainPart(attr("Target")!); + }); + } + zip.openReadStream(entry, (streamError, stream) => { + if (streamError || !stream) { + fail("ARCHIVE_INVALID"); + return; + } + active = stream; + let bytes = 0; + let decoder: TextDecoder | undefined; + stream.on("data", (chunk: Buffer) => { + if (settled) return; + bytes += chunk.byteLength; + total += chunk.byteLength; + if ( + bytes > limits.maxEntryBytes || + total > limits.maxArchiveBytes + ) { + fail("ZIP_BOMB"); + return; + } + if (!metadata) return; + try { + decoder ??= new TextDecoder( + chunk[0] === 0xff && chunk[1] === 0xfe + ? "utf-16le" + : chunk[0] === 0xfe && chunk[1] === 0xff + ? "utf-16be" + : "utf-8", + { fatal: true }, + ); + metadata.write(decoder.decode(chunk, { stream: true })); + } catch { + fail("ARCHIVE_INVALID"); + } + }); + stream.on("error", () => fail("ARCHIVE_INVALID")); + stream.on("end", () => { + active = undefined; + if (settled) return; + try { + metadata?.write(decoder?.decode() ?? "").close(); + } catch { + fail("ARCHIVE_INVALID"); + } + if (!settled) zip.readEntry(); + }); + }); + }); + zip.readEntry(); + }, + ); + }); +} + +async function extract(request: WorkerRequest): Promise { + const media = await inspectArchive(request.path, request); + // Loading a native binding is itself isolated and follows all refusal guards. + const require = createRequire(import.meta.url); + let anydoc: typeof import("@firecrawl/anydoc"); + try { + anydoc = require("@firecrawl/anydoc"); + } catch { + throw new ExtractionError("PARSER_UNAVAILABLE"); + } + const bytes = await readFile(request.path); + let text: string; + try { + text = await anydoc.toMarkdownBytes( + bytes, + "docx" as import("@firecrawl/anydoc").Format, + ); + } catch { + throw new ExtractionError("PARSE_FAILED"); + } + if (!text.trim()) throw new ExtractionError("NO_DOCUMENT_TEXT"); + if (Buffer.byteLength(text) > request.maxOutputBytes) + throw new ExtractionError("PARSER_OUTPUT_LIMIT"); + return { + profile: "document-v1", + text, + format: "markdown", + completeness: media ? "partial" : "complete", + may_omit_content: media, + pages_needing_ocr: [], + truncated: false, + parser: { name: "anydoc", version: "0.1.3" }, + }; +} + +if (import.meta.main) { + let serialized: string; + try { + const request = requestSchema.parse(JSON.parse(await Bun.stdin.text())); + serialized = JSON.stringify({ ok: true, result: await extract(request) }); + if (Buffer.byteLength(serialized) > request.maxOutputBytes) + throw new ExtractionError("PARSER_OUTPUT_LIMIT"); + } catch (error) { + serialized = JSON.stringify({ + ok: false, + code: error instanceof ExtractionError ? error.code : "PARSE_FAILED", + }); + } + await Bun.write(Bun.stdout, serialized); +} diff --git a/service/test/extraction.test.ts b/service/test/extraction.test.ts new file mode 100644 index 00000000..7ad6e3c9 --- /dev/null +++ b/service/test/extraction.test.ts @@ -0,0 +1,478 @@ +import { afterEach, beforeEach, expect, test } from "bun:test"; +import { SignJWT } from "jose"; +import { unzipSync, zipSync, strFromU8, strToU8 } from "fflate"; +import { mkdtemp, readdir, readFile, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { fileURLToPath } from "node:url"; +import { createApp } from "../src/app"; +import { configSchema, fromEnv } from "../src/config"; +import { DOCX_TYPE, resultSchema } from "../src/contract"; +import { runWorker, type Runner } from "../src/process"; + +const fixture = await Bun.file( + new URL("../../tests/fixtures/structured.docx", import.meta.url), +).bytes(); +const expectedText = + "# Quarterly Report\n\nThis document summarizes the results for the period.\n\n## Regional Totals\n\n| | | |\n| --- | --- | --- |\n| Region | Units | Revenue |\n| North | 1200 | 48000 |\n| South | 950 | 38000 |\n| East | 1430 | 57200 |\n\n**Totals are unaudited.**\n"; +const secret = "test-rag-key-with-more-than-32-characters"; +let tempRoot: string; +beforeEach(async () => { + tempRoot = await mkdtemp(join(tmpdir(), "rag-bun-test-")); +}); +afterEach(async () => { + await rm(tempRoot, { recursive: true, force: true }); +}); + +async function token( + options: { + scopes?: string[]; + key?: string; + audience?: string; + expiry?: number; + subject?: string; + } = {}, +) { + return new SignJWT({ scopes: options.scopes ?? ["rag:documents"] }) + .setProtectedHeader({ alg: "HS256" }) + .setIssuer("librechat") + .setAudience(options.audience ?? "rag-api") + .setSubject(options.subject ?? "owner") + .setExpirationTime(options.expiry ?? Math.floor(Date.now() / 1000) + 60) + .sign(new TextEncoder().encode(options.key ?? secret)); +} +function app(overrides: Parameters[0] = {}, runner?: Runner) { + return createApp({ enabled: true, secret, tempRoot, ...overrides }, runner); +} +function form( + bytes: Uint8Array = fixture, + mime = DOCX_TYPE, + name = "report.docx", + profile = "document-v1", +) { + const body = new FormData(); + body.append("profile", profile); + body.append("file", new File([new Uint8Array(bytes)], name, { type: mime })); + return body; +} +async function post( + application = app(), + body: BodyInit = form(), + jwt?: string, + signal?: AbortSignal, +) { + return application.request("/v1/extract", { + method: "POST", + headers: { Authorization: `Bearer ${jwt ?? (await token())}` }, + body, + signal, + }); +} +async function code(response: Response, status: number, error: string) { + expect(response.status).toBe(status); + expect(await response.json()).toEqual({ detail: { code: error } }); +} +async function clean() { + expect(await readdir(tempRoot)).toEqual([]); +} +function extraEntry(name: string, bytes: Uint8Array) { + return zipSync({ ...unzipSync(fixture), [name]: bytes }); +} + +// Same golden output as PR 330's Python prototype and Marco's pinned AnyDoc fixture. +test("real AnyDoc DOCX preserves the exact versioned Markdown contract", async () => { + const response = await post(); + expect(response.status).toBe(200); + expect(await response.json()).toEqual({ + profile: "document-v1", + text: expectedText, + format: "markdown", + completeness: "complete", + may_omit_content: false, + pages_needing_ocr: [], + truncated: false, + parser: { name: "anydoc", version: "0.1.3" }, + }); + await clean(); +}); +test("renamed DOCX follows MIME while generic DOCX follows filename", async () => { + for (const [mime, name] of [ + [DOCX_TYPE, "renamed.csv"], + ["application/octet-stream", "REPORT.DOCX"], + ] as const) { + expect((await post(app(), form(fixture, mime, name))).status).toBe(200); + } + await clean(); +}); +test("embedded image or object cannot claim complete inspection", async () => { + for (const name of ["word/media/scan.png", "word/embeddings/object.bin"]) { + const response = await post( + app(), + form(extraEntry(name, new Uint8Array([1, 2, 3]))), + ); + expect(response.status).toBe(200); + const result = resultSchema.parse(await response.json()); + expect(result.completeness).toBe("partial"); + expect(result.may_omit_content).toBe(true); + } + await clean(); +}); +test("image relationships and content types identify non-image filenames", async () => { + const data = unzipSync(fixture); + data["word/media/image.bin"] = new Uint8Array([1, 2, 3]); + data["[Content_Types].xml"] = strToU8( + strFromU8(data["[Content_Types].xml"]!).replace( + "", + '', + ), + ); + let response = await post(app(), form(zipSync(data))); + expect(response.status).toBe(200); + expect(resultSchema.parse(await response.json()).may_omit_content).toBe(true); + data["[Content_Types].xml"] = unzipSync(fixture)["[Content_Types].xml"]!; + data["word/_rels/document.xml.rels"] = strToU8( + '', + ); + response = await post(app(), form(zipSync(data))); + expect(response.status).toBe(200); + expect(resultSchema.parse(await response.json()).completeness).toBe( + "partial", + ); + await clean(); +}); +test("relationship-resolved main documents do not need a word directory", async () => { + const data = unzipSync(fixture); + for (const name of Object.keys(data)) { + if (!name.startsWith("word/")) continue; + data[name.replace("word/", "content/")] = data[name]!; + delete data[name]; + } + for (const name of ["_rels/.rels", "[Content_Types].xml"]) { + data[name] = strToU8( + strFromU8(data[name]!).replaceAll("word/", "content/"), + ); + } + const response = await post(app(), form(zipSync(data))); + expect(response.status).toBe(200); + const result = resultSchema.parse(await response.json()); + expect(result.text).toBe(expectedText); + await clean(); +}); +test("external main-part relationships and DTD metadata are hard refusals", async () => { + const data = unzipSync(fixture); + data["_rels/.rels"] = strToU8( + '', + ); + await code(await post(app(), form(zipSync(data))), 422, "ARCHIVE_INVALID"); + data["_rels/.rels"] = strToU8( + ']>&x;', + ); + await code(await post(app(), form(zipSync(data))), 422, "ARCHIVE_INVALID"); + await clean(); +}); +test("cover thumbnail does not cause paid OCR escalation", async () => { + const response = await post( + app(), + form(extraEntry("docProps/thumbnail.png", new Uint8Array([1]))), + ); + expect(response.status).toBe(200); + expect((await response.json()).may_omit_content).toBe(false); + await clean(); +}); +test("disabled service has no extraction route, even for a malformed body", async () => { + const response = await post(createApp(), "not multipart"); + await code(response, 404, "EXTRACTION_DISABLED"); + await clean(); + const health = await createApp().request("/health"); + expect(health.status).toBe(200); + expect((await health.json()).extraction_profiles).toEqual([]); +}); +test("enabled startup requires a dedicated service key and finite limits", () => { + expect(() => createApp({ enabled: true })).toThrow(); + for (const limits of [ + { timeoutMs: Infinity }, + { concurrent: 0 }, + { queued: -1 }, + { maxBodyBytes: 1 }, + ]) { + expect(() => + configSchema.parse({ enabled: true, secret, ...limits }), + ).toThrow(); + } + expect(() => fromEnv({ RAG_EXTRACTION_API_ENABLED: "yes" })).toThrow(); + expect(() => + fromEnv({ RAG_JWT_SECRET: secret, JWT_SECRET: secret }), + ).toThrow(); +}); +test("rejects absent, expired, wrong-audience, and session-key tokens before parsing", async () => { + const application = app(); + for (const jwt of [ + "", + await token({ expiry: 1 }), + await token({ audience: "session" }), + await token({ key: "different-session-signing-key-32-characters" }), + ]) { + await code( + await post(application, "not multipart", jwt), + 401, + "EXTRACTION_AUTH_REQUIRED", + ); + } + const noHeader = await application.request("/v1/extract", { method: "POST" }); + await code(noHeader, 401, "EXTRACTION_AUTH_REQUIRED"); + await clean(); +}); +test("inference scopes never grant document extraction", async () => { + for (const scopes of [[], ["rag:embed"], ["rag:rerank"]]) { + await code( + await post(app(), "not multipart", await token({ scopes })), + 403, + "EXTRACTION_FORBIDDEN", + ); + } + await clean(); +}); +test("legacy id token cannot reach the new service", async () => { + const legacy = await new SignJWT({ id: "owner" }) + .setProtectedHeader({ alg: "HS256" }) + .sign(new TextEncoder().encode(secret)); + await code(await post(app(), "bad", legacy), 401, "EXTRACTION_AUTH_REQUIRED"); + await clean(); +}); +test("unsupported type/profile and malformed multipart never reach native parsing", async () => { + await code( + await post(app(), form(fixture, "application/pdf")), + 415, + "UNSUPPORTED_DOCUMENT_TYPE", + ); + await code( + await post(app(), form(fixture, "text/markdown", "note.md")), + 415, + "UNSUPPORTED_DOCUMENT_TYPE", + ); + await code( + await post(app(), form(fixture, DOCX_TYPE, "report.docx", "raw-v1")), + 400, + "UNSUPPORTED_PROFILE", + ); + await code(await post(app(), "bad"), 400, "INVALID_MULTIPART"); + await clean(); +}); +test("duplicate fields, extra files and missing profile are rejected", async () => { + const duplicate = form(); + duplicate.append("profile", "document-v1"); + await code(await post(app(), duplicate), 400, "INVALID_MULTIPART"); + const files = form(); + files.append("file", new File([fixture], "second.docx", { type: DOCX_TYPE })); + await code(await post(app(), files), 400, "INVALID_MULTIPART"); + const missing = form(); + missing.delete("profile"); + await code(await post(app(), missing), 400, "INVALID_MULTIPART"); + await clean(); +}); +test("invalid archive errors never contain caller text", async () => { + await code( + await post(app(), form(new TextEncoder().encode("private contents"))), + 422, + "ARCHIVE_INVALID", + ); + await clean(); +}); +test("zip bomb, total size and entry counts are hard refusals", async () => { + const bomb = extraEntry("word/bomb.xml", new Uint8Array(20_000)); + await code( + await post(app({ maxEntryBytes: 10_000 }), form(bomb)), + 413, + "ZIP_BOMB", + ); + await code(await post(app({ maxArchiveBytes: 1000 })), 413, "ZIP_BOMB"); + await code(await post(app({ maxEntries: 2 })), 413, "ZIP_BOMB"); + await clean(); +}); +test("empty document returns no text instead of a successful extraction", async () => { + const xml = new TextEncoder().encode( + '', + ); + await code( + await post(app(), form(extraEntry("word/document.xml", xml))), + 422, + "NO_DOCUMENT_TEXT", + ); + await clean(); +}); +test("file size and serialized native output have independent ceilings", async () => { + await code(await post(app({ maxFileBytes: 500 })), 413, "PARSER_INPUT_LIMIT"); + await code( + await post(app({ maxOutputBytes: 100 })), + 413, + "PARSER_OUTPUT_LIMIT", + ); + await clean(); +}); +test("parent caps IPC even if a child ignores its output limit", async () => { + const runner: Runner = (request, signal) => + runWorker(request, signal, [ + process.execPath, + "-e", + 'process.stdout.write("x".repeat(1000))', + ]); + await code( + await post(app({ maxOutputBytes: 100 }, runner)), + 413, + "PARSER_OUTPUT_LIMIT", + ); + await clean(); +}); +test("native children do not reload secrets from dotenv files", async () => { + await Bun.write( + join(tempRoot, ".env"), + "NATIVE_ENV_CANARY=should-not-reach-parser\n", + ); + const result = { + profile: "document-v1", + text: "isolated", + format: "markdown", + completeness: "complete", + may_omit_content: false, + pages_needing_ocr: [], + truncated: false, + parser: { name: "anydoc", version: "0.1.3" }, + }; + const script = `const result = ${JSON.stringify(result)}; if (process.env.NATIVE_ENV_CANARY) result.text = "leaked"; console.log(JSON.stringify({ok:true,result}));`; + const response = await runWorker( + { + path: "unused", + maxOutputBytes: 4096, + maxEntryBytes: 4096, + maxArchiveBytes: 4096, + maxEntries: 5, + }, + new AbortController().signal, + [process.execPath, "--cwd", tempRoot, "-e", script], + ); + expect(response.text).toBe("isolated"); +}); + +test("child crash and malformed response are sanitized and permit retry", async () => { + for (const source of [ + "process.exit(11)", + 'console.log("private native failure")', + 'console.log("null")', + ]) { + const runner: Runner = (request, signal) => + runWorker(request, signal, [process.execPath, "-e", source]); + await code(await post(app({}, runner)), 503, "PARSER_CRASH"); + await clean(); + } + expect((await post()).status).toBe(200); + await clean(); +}); + +const sleepPath = fileURLToPath(new URL("./sleep.fixture.ts", import.meta.url)); +const sleeping: Runner = (request, signal) => + runWorker(request, signal, [process.execPath, sleepPath, request.path]); +async function childPid() { + for (let i = 0; i < 200; i++) { + for (const dir of await readdir(tempRoot)) { + try { + return Number( + await readFile(join(tempRoot, dir, "input.docx.pid"), "utf8"), + ); + } catch { + /* Not started yet. */ + } + } + await Bun.sleep(10); + } + throw new Error("Child never started"); +} +function reaped(pid: number) { + expect(() => process.kill(pid, 0)).toThrow(); +} +test("overall deadline kills and reaps a running native process before cleanup", async () => { + const pending = post(app({ timeoutMs: 400 }, sleeping)); + const pid = await childPid(); + await code(await pending, 504, "PARSER_TIMEOUT"); + reaped(pid); + await clean(); +}); +test("abort kills and reaps the child, cleans temp files and permits the next request", async () => { + let calls = 0; + const runner: Runner = (request, signal) => + ++calls === 1 ? sleeping(request, signal) : runWorker(request, signal); + const application = app({}, runner); + const controller = new AbortController(); + const pending = post(application, form(), await token(), controller.signal); + const pid = await childPid(); + controller.abort(); + await code(await pending, 408, "REQUEST_CANCELLED"); + reaped(pid); + await clean(); + expect((await post(application)).status).toBe(200); + await clean(); +}); +test("overload refuses before reading or staging another request body", async () => { + const application = app({ concurrent: 1, queued: 0 }, sleeping); + const controller = new AbortController(); + const pending = post(application, form(), await token(), controller.signal); + const pid = await childPid(); + let pulls = 0; + const body = new ReadableStream( + { + pull(control) { + pulls++; + control.enqueue(new Uint8Array([1])); + }, + }, + { highWaterMark: 0 }, + ); + const request = new Request("http://test/v1/extract", { + method: "POST", + headers: { + Authorization: `Bearer ${await token()}`, + "Content-Type": "multipart/form-data; boundary=test", + }, + body, + }); + await code(await application.fetch(request), 429, "CONCURRENCY_LIMIT"); + expect(pulls).toBe(0); + controller.abort(); + await pending; + reaped(pid); + await clean(); +}); +test("body limit is counted for chunked requests, not trusted Content-Length", async () => { + const application = app({ maxFileBytes: 100, maxBodyBytes: 200 }); + let pulls = 0; + let cancelled = false; + const prefix = + '--test\r\nContent-Disposition: form-data; name="file"; filename="r.docx"\r\nContent-Type: ' + + DOCX_TYPE + + "\r\n\r\n"; + const body = new ReadableStream( + { + pull(control) { + pulls++; + control.enqueue( + new TextEncoder().encode(pulls === 1 ? prefix : "x".repeat(256)), + ); + }, + cancel() { + cancelled = true; + }, + }, + { highWaterMark: 0 }, + ); + const request = new Request("http://test/v1/extract", { + method: "POST", + headers: { + Authorization: `Bearer ${await token()}`, + "Content-Type": "multipart/form-data; boundary=test", + }, + body, + }); + await code(await application.fetch(request), 413, "PARSER_INPUT_LIMIT"); + expect(pulls).toBeLessThan(5); + expect(cancelled).toBe(true); + await clean(); +}); diff --git a/service/test/image.sh b/service/test/image.sh new file mode 100644 index 00000000..fd67ca56 --- /dev/null +++ b/service/test/image.sh @@ -0,0 +1,28 @@ +#!/usr/bin/env bash +set -euo pipefail +image=${1:?Pass the locally built image name} +root=$(git rev-parse --show-toplevel) +docker run --rm --entrypoint sh "$image" -c '! command -v python && ! command -v python3' +docker run --rm -i --entrypoint bun "$image" -e ' +import { createApp } from "./src/app.ts"; +import { SignJWT } from "jose"; +const secret = "image-smoke-test-key-not-a-real-service-key"; +const data = await Bun.stdin.bytes(); +const application = createApp({enabled: true, secret}); +const server = Bun.serve({hostname: "127.0.0.1", port: 0, fetch: application.fetch}); +try { + const health = await fetch(`http://127.0.0.1:${server.port}/health`); + if (health.status !== 200) throw new Error("Health failed"); + const token = await new SignJWT({scopes: ["rag:documents"]}).setProtectedHeader({alg: "HS256"}) + .setIssuer("librechat").setAudience("rag-api").setSubject("owner").setExpirationTime("1m") + .sign(new TextEncoder().encode(secret)); + const body = new FormData(); body.append("profile", "document-v1"); + body.append("file", new File([data], "report.docx", {type: "application/vnd.openxmlformats-officedocument.wordprocessingml.document"})); + const response = await fetch(`http://127.0.0.1:${server.port}/v1/extract`, {method: "POST", body, headers: {Authorization: `Bearer ${token}`}}); + const result = await response.json(); + if (response.status !== 200 || !result.text?.includes("Regional Totals") || result.parser?.name !== "anydoc") { + throw new Error(`Native DOCX failed: ${response.status}`); + } + console.log("PASS: Bun listener, health and native DOCX extraction in a Python-free image"); +} finally { await server.stop(true); } +' < "$root/tests/fixtures/structured.docx" diff --git a/service/test/lifecycle.test.ts b/service/test/lifecycle.test.ts new file mode 100644 index 00000000..577c87b1 --- /dev/null +++ b/service/test/lifecycle.test.ts @@ -0,0 +1,154 @@ +import { expect, test } from "bun:test"; +import { SignJWT } from "jose"; +import { mkdtemp, readdir, readFile, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { fileURLToPath } from "node:url"; +import { Admission } from "../src/admission"; +import { createApp } from "../src/app"; +import { configSchema } from "../src/config"; +import { listen } from "../src/server"; +import { DOCX_TYPE } from "../src/contract"; +import { runWorker, type Runner } from "../src/process"; + +const fixture = await Bun.file( + new URL("../../tests/fixtures/structured.docx", import.meta.url), +).bytes(); +const sleepPath = fileURLToPath(new URL("./sleep.fixture.ts", import.meta.url)); +const secret = "lifecycle-test-service-signing-key-32-characters"; + +async function headers() { + const token = await new SignJWT({ scopes: ["rag:documents"] }) + .setProtectedHeader({ alg: "HS256" }) + .setIssuer("librechat") + .setAudience("rag-api") + .setSubject("owner") + .setExpirationTime("1m") + .sign(new TextEncoder().encode(secret)); + return { Authorization: `Bearer ${token}` }; +} +function form() { + const body = new FormData(); + body.append("profile", "document-v1"); + body.append("file", new File([fixture], "report.docx", { type: DOCX_TYPE })); + return body; +} +async function pid(temp: string) { + for (let i = 0; i < 200; i++) { + for (const dir of await readdir(temp)) { + try { + return Number( + await readFile(join(temp, dir, "input.docx.pid"), "utf8"), + ); + } catch { + /* Child is not running yet. */ + } + } + await Bun.sleep(10); + } + throw new Error("Child never started"); +} + +test("cancelled queued work never consumes a slot and remaining work stays FIFO", async () => { + const admission = new Admission(1, 2); + const first = await admission.acquire(new AbortController().signal); + const cancelled = new AbortController(); + const second = admission.acquire(cancelled.signal); + let granted = false; + const third = admission + .acquire(new AbortController().signal) + .then((release) => { + granted = true; + return release; + }); + cancelled.abort(); + await expect(second).rejects.toMatchObject({ code: "REQUEST_CANCELLED" }); + expect(granted).toBe(false); + first(); + const release = await third; + expect(granted).toBe(true); + release(); + release(); + const next = await admission.acquire(new AbortController().signal); + next(); +}); + +test("listener waits beyond ten seconds for the configured extraction deadline", async () => { + const tempRoot = await mkdtemp(join(tmpdir(), "rag-long-parse-")); + const config = configSchema.parse({ + enabled: true, + secret, + tempRoot, + timeoutMs: 14_000, + }); + const runner: Runner = async (request, signal) => { + await Bun.sleep(11_000); + return runWorker(request, signal); + }; + const server = listen(createApp(config, runner), config, 0, "127.0.0.1"); + try { + const response = await fetch(`http://127.0.0.1:${server.port}/v1/extract`, { + method: "POST", + headers: await headers(), + body: form(), + }); + expect(response.status).toBe(200); + expect((await response.json()).text).toContain("Quarterly Report"); + } finally { + await server.stop(true); + await rm(tempRoot, { recursive: true, force: true }); + } +}, 20_000); + +test("actual HTTP disconnect kills the native child and cleans up before retry", async () => { + const tempRoot = await mkdtemp(join(tmpdir(), "rag-disconnect-")); + let calls = 0; + const runner: Runner = (request, signal) => + ++calls === 1 + ? runWorker(request, signal, [process.execPath, sleepPath, request.path]) + : runWorker(request, signal); + const application = createApp( + { + enabled: true, + secret, + tempRoot, + concurrent: 1, + queued: 0, + timeoutMs: 3000, + }, + runner, + ); + const server = Bun.serve({ + hostname: "127.0.0.1", + port: 0, + fetch: application.fetch, + }); + const controller = new AbortController(); + try { + const pending = fetch(`http://127.0.0.1:${server.port}/v1/extract`, { + method: "POST", + headers: await headers(), + body: form(), + signal: controller.signal, + }); + void pending.catch(() => {}); + const child = await pid(tempRoot); + controller.abort(); + await expect(pending).rejects.toThrow(); + for (let i = 0; i < 200 && (await readdir(tempRoot)).length; i++) + await Bun.sleep(10); + expect(await readdir(tempRoot)).toEqual([]); + expect(() => process.kill(child, 0)).toThrow(); + const retry = await fetch(`http://127.0.0.1:${server.port}/v1/extract`, { + method: "POST", + headers: await headers(), + body: form(), + }); + expect(retry.status).toBe(200); + expect((await retry.json()).text).toContain("Quarterly Report"); + } finally { + controller.abort(); + await server.stop(true); + await rm(tempRoot, { recursive: true, force: true }); + } +}); diff --git a/service/test/sleep.fixture.ts b/service/test/sleep.fixture.ts new file mode 100644 index 00000000..25fc93a5 --- /dev/null +++ b/service/test/sleep.fixture.ts @@ -0,0 +1,6 @@ +export {}; + +const path = process.argv[2]; +if (!path) throw new Error("Missing fixture path"); +await Bun.write(`${path}.pid`, String(process.pid)); +await Bun.sleep(60_000); diff --git a/service/tsconfig.json b/service/tsconfig.json new file mode 100644 index 00000000..c0d4f1a1 --- /dev/null +++ b/service/tsconfig.json @@ -0,0 +1,15 @@ +{ + "compilerOptions": { + "target": "ES2023", + "lib": ["ES2023", "DOM"], + "module": "ESNext", + "moduleResolution": "Bundler", + "types": ["bun"], + "strict": true, + "noEmit": true, + "noUncheckedIndexedAccess": true, + "allowImportingTsExtensions": true, + "skipLibCheck": true + }, + "include": ["src/**/*.ts", "test/**/*.ts"] +} diff --git a/tests/fixtures/structured.docx b/tests/fixtures/structured.docx new file mode 100644 index 00000000..eff663f8 Binary files /dev/null and b/tests/fixtures/structured.docx differ