From c3ce95fa973f22b5dc2ad166dc9f89550e90d308 Mon Sep 17 00:00:00 2001 From: Lia Date: Sun, 4 Oct 2026 18:06:10 +0000 Subject: [PATCH 1/3] =?UTF-8?q?=F0=9F=A7=AA=20feat:=20Add=20Scoped=20Bun?= =?UTF-8?q?=20RAG=20v2=20PoC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .github/workflows/v2.yml | 23 ++ v2/.dockerignore | 5 + v2/.gitignore | 3 + v2/Dockerfile | 16 + v2/bun.lock | 81 +++++ v2/compose.yaml | 43 +++ v2/package.json | 30 ++ v2/scripts/benchmark.ts | 216 ++++++++++++++ v2/scripts/qualify.sh | 26 ++ v2/src/app.ts | 297 +++++++++++++++++++ v2/src/auth.ts | 135 +++++++++ v2/src/chunking.ts | 67 +++++ v2/src/clickhouse.ts | 572 ++++++++++++++++++++++++++++++++++++ v2/src/config.ts | 59 ++++ v2/src/contracts.ts | 194 ++++++++++++ v2/src/embeddings.ts | 223 ++++++++++++++ v2/src/extraction.ts | 85 ++++++ v2/src/limits.ts | 85 ++++++ v2/src/main.ts | 77 +++++ v2/src/migrate.ts | 32 ++ v2/src/pipeline.ts | 275 +++++++++++++++++ v2/src/schema.ts | 45 +++ v2/test/app.test.ts | 206 +++++++++++++ v2/test/auth.test.ts | 68 +++++ v2/test/chunking.test.ts | 77 +++++ v2/test/clickhouse.test.ts | 44 +++ v2/test/embeddings.test.ts | 105 +++++++ v2/test/extraction.test.ts | 59 ++++ v2/test/helpers.ts | 144 +++++++++ v2/test/integration.test.ts | 242 +++++++++++++++ v2/test/pipeline.test.ts | 261 ++++++++++++++++ v2/tsconfig.json | 16 + 32 files changed, 3811 insertions(+) create mode 100644 .github/workflows/v2.yml create mode 100644 v2/.dockerignore create mode 100644 v2/.gitignore create mode 100644 v2/Dockerfile create mode 100644 v2/bun.lock create mode 100644 v2/compose.yaml create mode 100644 v2/package.json create mode 100644 v2/scripts/benchmark.ts create mode 100755 v2/scripts/qualify.sh create mode 100644 v2/src/app.ts create mode 100644 v2/src/auth.ts create mode 100644 v2/src/chunking.ts create mode 100644 v2/src/clickhouse.ts create mode 100644 v2/src/config.ts create mode 100644 v2/src/contracts.ts create mode 100644 v2/src/embeddings.ts create mode 100644 v2/src/extraction.ts create mode 100644 v2/src/limits.ts create mode 100644 v2/src/main.ts create mode 100644 v2/src/migrate.ts create mode 100644 v2/src/pipeline.ts create mode 100644 v2/src/schema.ts create mode 100644 v2/test/app.test.ts create mode 100644 v2/test/auth.test.ts create mode 100644 v2/test/chunking.test.ts create mode 100644 v2/test/clickhouse.test.ts create mode 100644 v2/test/embeddings.test.ts create mode 100644 v2/test/extraction.test.ts create mode 100644 v2/test/helpers.ts create mode 100644 v2/test/integration.test.ts create mode 100644 v2/test/pipeline.test.ts create mode 100644 v2/tsconfig.json diff --git a/.github/workflows/v2.yml b/.github/workflows/v2.yml new file mode 100644 index 00000000..616410d5 --- /dev/null +++ b/.github/workflows/v2.yml @@ -0,0 +1,23 @@ +name: RAG v2 PoC +on: + pull_request: + paths: ['v2/**', '.github/workflows/v2.yml'] + push: + branches: [main] + paths: ['v2/**', '.github/workflows/v2.yml'] +permissions: + contents: read +jobs: + focused: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: oven-sh/setup-bun@v2 + with: + bun-version: 1.4.2 + - run: bun install --frozen-lockfile --ignore-scripts + working-directory: v2 + - run: bun run typecheck && bun run format:check && bun test test && bun run build + working-directory: v2 + - run: bun run test:container + working-directory: v2 diff --git a/v2/.dockerignore b/v2/.dockerignore new file mode 100644 index 00000000..322853d7 --- /dev/null +++ b/v2/.dockerignore @@ -0,0 +1,5 @@ +node_modules +dist +test +scripts +.qualification diff --git a/v2/.gitignore b/v2/.gitignore new file mode 100644 index 00000000..f82270fa --- /dev/null +++ b/v2/.gitignore @@ -0,0 +1,3 @@ +node_modules/ +dist/ +.qualification/ diff --git a/v2/Dockerfile b/v2/Dockerfile new file mode 100644 index 00000000..b8a9f38c --- /dev/null +++ b/v2/Dockerfile @@ -0,0 +1,16 @@ +FROM oven/bun:1.4.2-debian AS build +WORKDIR /app +COPY package.json bun.lock ./ +RUN bun install --frozen-lockfile --ignore-scripts +COPY tsconfig.json ./ +COPY src ./src +RUN bun run typecheck && bun run build +FROM oven/bun:1.4.2-debian +WORKDIR /app +COPY --from=build /app/node_modules ./node_modules +COPY --from=build /app/dist ./dist +COPY package.json ./ +USER bun +ENV RAG_V2_HOST=0.0.0.0 +EXPOSE 8001 +CMD ["bun", "dist/main.js"] diff --git a/v2/bun.lock b/v2/bun.lock new file mode 100644 index 00000000..b5b428d0 --- /dev/null +++ b/v2/bun.lock @@ -0,0 +1,81 @@ +{ + "lockfileVersion": 2, + "configVersion": 1, + "workspaces": { + "": { + "name": "@librechat/rag-api-v2-poc", + "dependencies": { + "@clickhouse/client": "1.23.1", + "hono": "4.13.13", + "jose": "6.2.12", + "zod": "4.6.5", + }, + "devDependencies": { + "@types/bun": "1.4.2", + "prettier": "3.9.9", + "typescript": "7.0.2", + }, + }, + }, + "packages": { + "@clickhouse/client": ["@clickhouse/client@1.23.1", "", {}, "sha512-vs3/Zc1dHvT171btW5nMoPsPCJ6QVJ5pp7obxzO5sjqwFx/jjz9wwCAqcFOdc2DhprugDBaVn+4dVY8hG3A9nw=="], + + "@types/bun": ["@types/bun@1.4.2", "", { "dependencies": { "bun-types": "1.4.2" } }, "sha512-GimotNn7+ZV0uVArItBbriZsR1oNf0+WTzPkdcFrzShI7k2norL0uzEaJT8T33dWr7O/c9ZDuAFQrctKCi72oQ=="], + + "@types/node": ["@types/node@26.6.4", "", { "dependencies": { "undici-types": "~8.9.0" } }, "sha512-ldVPDCzj7fsaGZrLB0NuHuTvJcsNasysBAqMolr/cgxrLd1xbqxIr3XJiPnHHJUCxj5sNF1vnRj9aWnrVh5Jcg=="], + + "@typescript/typescript-aix-ppc64": ["@typescript/typescript-aix-ppc64@7.0.2", "", { "os": "aix", "cpu": "ppc64" }, "sha512-MTKKkWB7p/0E9xi1d1tHtZ5PiLkGEMIq88pK2CubZjOsLtYTLqhgIgi6zepFa+9GHZ6h05NMCkQxGKiPXMxXtQ=="], + + "@typescript/typescript-darwin-arm64": ["@typescript/typescript-darwin-arm64@7.0.2", "", { "os": "darwin", "cpu": "arm64" }, "sha512-gowzar9MwS/aRWp6f3a4KUqzRjAZjOsmGNCM6LcTgXum+dBfgsBVMN+AgvOCCbguXyick6LJhpBszxMebJ8syA=="], + + "@typescript/typescript-darwin-x64": ["@typescript/typescript-darwin-x64@7.0.2", "", { "os": "darwin", "cpu": "x64" }, "sha512-SZ9xZInqApNlNGc9s0W1VSsktYSOe9cFqNOIqmN1Gs8SmkjKZYFt017G4VwPxASInODuAdbTW7sXiFUf893RgA=="], + + "@typescript/typescript-freebsd-arm64": ["@typescript/typescript-freebsd-arm64@7.0.2", "", { "os": "freebsd", "cpu": "arm64" }, "sha512-W5NH4y/J0plIIS5b2xvTEkU7JFxyqdMAOgf+Ilhl0vHQXKO5dZoxd+C/jEtq56c4F3wk71RB4BMRQ2XdI+bwYQ=="], + + "@typescript/typescript-freebsd-x64": ["@typescript/typescript-freebsd-x64@7.0.2", "", { "os": "freebsd", "cpu": "x64" }, "sha512-UMGDx5sTpzNw3WiPebH7l90IWfJggEd+egHt/q6p7/Cm3zqoV7VxkGXt+3DxPIw8CcmvAB0j3sVVfbhX+M4Tpw=="], + + "@typescript/typescript-linux-arm": ["@typescript/typescript-linux-arm@7.0.2", "", { "os": "linux", "cpu": "arm" }, "sha512-gffT3xPz9sR7j/YJExkyPntrI0P2EP9XbOyWzth2/Gs0RstK+90RBcO0ncXoXy/beYll1SXw846Nf2zdnEz0QQ=="], + + "@typescript/typescript-linux-arm64": ["@typescript/typescript-linux-arm64@7.0.2", "", { "os": "linux", "cpu": "arm64" }, "sha512-Qh4eU4/y3yDjnfjjyPYihMj5/ODIlmt+Bzu17OI+fiSRDW57QmU5SiN63exPRNJPKUzcc1INa1NXdrJ+MqHjUQ=="], + + "@typescript/typescript-linux-loong64": ["@typescript/typescript-linux-loong64@7.0.2", "", { "os": "linux", "cpu": "none" }, "sha512-uEHck9i8hoAzXPiYRib1O7miOnz23SxIeVl6F4LXox+qov1K35jHcEW6VHKvZI+pyvl7fZEP4MCU5LYvIq1GuQ=="], + + "@typescript/typescript-linux-mips64el": ["@typescript/typescript-linux-mips64el@7.0.2", "", { "os": "linux", "cpu": "none" }, "sha512-R4KvAMnE43W5Qeqb0Ly56O3mWMWIAgsMyz36DCaycd5nbg/9kzm0liw3JocfRqyJY0KPmzFjbswozXyW0DnIYA=="], + + "@typescript/typescript-linux-ppc64": ["@typescript/typescript-linux-ppc64@7.0.2", "", { "os": "linux", "cpu": "ppc64" }, "sha512-DORx5b3sd/4S7eayxm4FQv+A7CrkUIGRaHiwI8oiHTAI1fAPWhF4J0vAlkC8biAlHSVVwxMQ3tjZ2/DVbnQiiA=="], + + "@typescript/typescript-linux-riscv64": ["@typescript/typescript-linux-riscv64@7.0.2", "", { "os": "linux", "cpu": "none" }, "sha512-wf0jqEDOjrPRnKwYRyyJDRo11KMbvMFrU+q4zqKyChODBzvlkbhNQfKvLxQCcwTpdDaXSHZTVuh0JoCrKCUMHQ=="], + + "@typescript/typescript-linux-s390x": ["@typescript/typescript-linux-s390x@7.0.2", "", { "os": "linux", "cpu": "s390x" }, "sha512-IkwJc3L7yhytWd/ewjyxNDfOmswCm9GWMJT/ue/dU4aZNbwZeYAetq42VyLmsmSjvoX7z74X6ZaYCtzAr0EuGw=="], + + "@typescript/typescript-linux-x64": ["@typescript/typescript-linux-x64@7.0.2", "", { "os": "linux", "cpu": "x64" }, "sha512-EYdf2cNg7rgCWJnxCdJ+F3V39O8ihb37eHAu1LK8oAFizgTQbPOK7zHHXbPt8rX24COqODXeI3sIf0fCXG7H/A=="], + + "@typescript/typescript-netbsd-arm64": ["@typescript/typescript-netbsd-arm64@7.0.2", "", { "os": "none", "cpu": "arm64" }, "sha512-+polYF4MF04aPpO5FTkHran9yUQDSXqy5GiSDKpsll5jy3l3+g9QLhpf39T+ePtefhXLOGrLl0QIjkQP6VnelA=="], + + "@typescript/typescript-netbsd-x64": ["@typescript/typescript-netbsd-x64@7.0.2", "", { "os": "none", "cpu": "x64" }, "sha512-8YIT0EHM/3dq10ZOVF/A7pc/YSMtbcecct4rWtexrnSCHOPcpC2KTLXfTCR6vDpnSiY12heNb1GiN/wu+T/FyA=="], + + "@typescript/typescript-openbsd-arm64": ["@typescript/typescript-openbsd-arm64@7.0.2", "", { "os": "openbsd", "cpu": "arm64" }, "sha512-APT8+ClYnuYm1u9+kgGXoMj2VzWzcymwh2gNSQVySHfkRDGOTVkoWLjCmOQSaO+PoqQ57B0flRp9SA+7GnnkzQ=="], + + "@typescript/typescript-openbsd-x64": ["@typescript/typescript-openbsd-x64@7.0.2", "", { "os": "openbsd", "cpu": "x64" }, "sha512-yX7s+Q0Dln0Dt9tEzZsAjXXR/+ytBM7AlglaqyeMPxQszJ1JhlJdZ6jLA+IzldHtflX81em7lDao1xXu+aRRkg=="], + + "@typescript/typescript-sunos-x64": ["@typescript/typescript-sunos-x64@7.0.2", "", { "os": "sunos", "cpu": "x64" }, "sha512-dLJDGaLZ1D4HPQn62u1n8mBDkJREwMsAkCdkwd4Ieqw+x3TUyTsqY0YiBCtE6H6OzzgGk3iuZ3vFWRS+E8/d1g=="], + + "@typescript/typescript-win32-arm64": ["@typescript/typescript-win32-arm64@7.0.2", "", { "os": "win32", "cpu": "arm64" }, "sha512-Gyl1Vy6OsWesLzmq+EP0Fb7b4Nid5232AvcA2SFcdYreldpNtYFFofPjnt62y9hQy7VTaZp65ICJjuAQRaVcIQ=="], + + "@typescript/typescript-win32-x64": ["@typescript/typescript-win32-x64@7.0.2", "", { "os": "win32", "cpu": "x64" }, "sha512-0BQ3HkAHHlKLSp1qRvf3SUhGpGsDuhB/jgFw75guyqbxJqEaS0Cw/VFO8i2nHglJUzQCRtMMR/IBAKE3ETMC4g=="], + + "bun-types": ["bun-types@1.4.2", "", { "dependencies": { "@types/node": "*" } }, "sha512-bxV1FgK7yBIzjRe5zBozIM4Bem11ZJcCXSrjWRG3YWLt8yFDePu4cLjpebO8OvPeIE9trbyPF4fuj3Cia4Fj3w=="], + + "hono": ["hono@4.13.13", "", {}, "sha512-CQ46U0ZkAGmbT/4UxdzzGJpacP2IeKgY4a5/tOI9AABbpOMfK739wfDXmv1usCk+3RkKj1hQy4/fjhiwa2xlrA=="], + + "jose": ["jose@6.2.12", "", {}, "sha512-9NiFmJEex0sy2Dk58j2UGBSHgUs2ypF9eZSu4L6vjOX3Dp96Sw1F3uL+H+D1sx02jZZdzUT0HgvCy59CuvXcWw=="], + + "prettier": ["prettier@3.9.9", "", { "bin": { "prettier": "bin/prettier.cjs" } }, "sha512-Z/CJHIkdujO/OtN7nXUii0Rf3VT5SRuhjBA82Xvu2XhBUgX3nhP67T0LHceBdQLex7OOFGTox+Q5Yg8Jk2Qivg=="], + + "typescript": ["typescript@7.0.2", "", { "optionalDependencies": { "@typescript/typescript-aix-ppc64": "7.0.2", "@typescript/typescript-darwin-arm64": "7.0.2", "@typescript/typescript-darwin-x64": "7.0.2", "@typescript/typescript-freebsd-arm64": "7.0.2", "@typescript/typescript-freebsd-x64": "7.0.2", "@typescript/typescript-linux-arm": "7.0.2", "@typescript/typescript-linux-arm64": "7.0.2", "@typescript/typescript-linux-loong64": "7.0.2", "@typescript/typescript-linux-mips64el": "7.0.2", "@typescript/typescript-linux-ppc64": "7.0.2", "@typescript/typescript-linux-riscv64": "7.0.2", "@typescript/typescript-linux-s390x": "7.0.2", "@typescript/typescript-linux-x64": "7.0.2", "@typescript/typescript-netbsd-arm64": "7.0.2", "@typescript/typescript-netbsd-x64": "7.0.2", "@typescript/typescript-openbsd-arm64": "7.0.2", "@typescript/typescript-openbsd-x64": "7.0.2", "@typescript/typescript-sunos-x64": "7.0.2", "@typescript/typescript-win32-arm64": "7.0.2", "@typescript/typescript-win32-x64": "7.0.2" }, "bin": { "tsc": "bin/tsc" } }, "sha512-8FYau96o3NKOhbjKi/qNvG/W5jhzxkbdm5sj9AbZ/5T5sWqn3hJgLfGx27sRKZWTvyzCP8dLRBTf5tBTSRVUNA=="], + + "undici-types": ["undici-types@8.9.0", "", {}, "sha512-KTDyRTYX8sWmKXAikPHHSyc63CRPETMctyjKFupcC6OBLXT3xsN0e9aF7m+mIXutFWpUXuedtowG7iLOzp0kQg=="], + + "zod": ["zod@4.6.5", "", {}, "sha512-v5l/aFXZQeai4awLbOpSoHecE9UiMrnfx75tEXLjNonXVARxQ5mOeipTjROUchszUNCqnE+hqAMujRsRHsut2Q=="], + } +} diff --git a/v2/compose.yaml b/v2/compose.yaml new file mode 100644 index 00000000..f8c0b573 --- /dev/null +++ b/v2/compose.yaml @@ -0,0 +1,43 @@ +services: + clickhouse: + image: clickhouse/clickhouse-server:26.8 + environment: + CLICKHOUSE_DB: rag_v2 + CLICKHOUSE_USER: rag + CLICKHOUSE_PASSWORD: ${RAG_V2_CLICKHOUSE_PASSWORD:?Set a database password} + volumes: + - clickhouse-data:/var/lib/clickhouse + networks: [database] + healthcheck: + test: [CMD, clickhouse-client, --query, SELECT 1] + interval: 5s + timeout: 3s + retries: 20 + rag: + build: . + init: true + read_only: true + cap_drop: [ALL] + security_opt: [no-new-privileges:true] + ports: ["127.0.0.1:8001:8001"] + environment: + RAG_V2_JWKS_JSON: ${RAG_V2_JWKS_JSON:?Set public verification keys} + RAG_V2_JWT_ISSUER: ${RAG_V2_JWT_ISSUER:-librechat} + RAG_V2_CLICKHOUSE_URL: http://clickhouse:8123 + RAG_V2_ALLOW_INSECURE_CLICKHOUSE: "true" + RAG_V2_CLICKHOUSE_USER: rag + RAG_V2_CLICKHOUSE_PASSWORD: ${RAG_V2_CLICKHOUSE_PASSWORD:?Set a database password} + RAG_V2_EMBEDDING_API_KEY: ${RAG_V2_EMBEDDING_API_KEY:?Set an embedding key} + RAG_V2_MODE: coordinator + RAG_V2_SINGLE_WRITER: "true" + command: [sh, -c, "bun dist/migrate.js && exec bun dist/main.js"] + depends_on: + clickhouse: + condition: service_healthy + networks: [database, outbound] +networks: + database: + internal: true + outbound: {} +volumes: + clickhouse-data: {} diff --git a/v2/package.json b/v2/package.json new file mode 100644 index 00000000..23d204b6 --- /dev/null +++ b/v2/package.json @@ -0,0 +1,30 @@ +{ + "name": "@librechat/rag-api-v2-poc", + "version": "0.1.0", + "private": true, + "type": "module", + "engines": { + "bun": "1.4.2" + }, + "scripts": { + "start": "bun src/main.ts", + "migrate": "bun src/migrate.ts", + "test": "bun test test", + "typecheck": "tsc --noEmit", + "build": "bun build src/main.ts src/migrate.ts --target=bun --packages=external --outdir=dist", + "format:check": "prettier --check src test scripts package.json tsconfig.json compose.yaml", + "test:container": "bash scripts/qualify.sh", + "bench": "bun scripts/benchmark.ts" + }, + "dependencies": { + "@clickhouse/client": "1.23.1", + "hono": "4.13.13", + "jose": "6.2.12", + "zod": "4.6.5" + }, + "devDependencies": { + "@types/bun": "1.4.2", + "prettier": "3.9.9", + "typescript": "7.0.2" + } +} diff --git a/v2/scripts/benchmark.ts b/v2/scripts/benchmark.ts new file mode 100644 index 00000000..2a7557b2 --- /dev/null +++ b/v2/scripts/benchmark.ts @@ -0,0 +1,216 @@ +import type { ClickHouseSettings } from "@clickhouse/client"; +import { + databaseClient, + ClickHouseStore, + liveWhere, + vectorSQL, +} from "../src/clickhouse"; +import { migrate } from "../src/schema"; +import { createApp } from "../src/app"; +import { auth, token } from "../test/helpers"; +import type { Document, StoredChunk } from "../src/contracts"; +import type { EmbeddingProvider } from "../src/embeddings"; + +const url = process.env.RAG_V2_TEST_CLICKHOUSE_URL; +if (!url) throw Error("RAG_V2_TEST_CLICKHOUSE_URL_REQUIRED"); +const dimensions = Number(process.env.RAG_V2_BENCH_DIMENSIONS ?? 1536); +const count = Number(process.env.RAG_V2_BENCH_ROWS ?? 10000); +if (!Number.isInteger(count) || count < 100 || count > 100000) + throw Error("INVALID_BENCH_ROWS"); +const database = `ragv2_bench_${crypto.randomUUID().replaceAll("-", "")}`; +const config = { url, database: "default", username: "default", password: "" }; +const admin = databaseClient(config); +const client = databaseClient({ ...config, database }); +const store = new ClickHouseStore(client, database, dimensions); +const signal = () => AbortSignal.timeout(30000); +let seed = 17; +const vector = () => + Array.from({ length: dimensions }, () => { + seed = (Math.imul(seed, 1664525) + 1013904223) >>> 0; + return seed / 4294967296 - 0.5; + }); +const reference = vector(); +const scope = { tenantId: "benchmark", namespaceId: "library-a" }; +const generation = "a".repeat(32); +const spaceId = `synthetic-${dimensions}`; +const provider: EmbeddingProvider = { + spaceId, + dimensions, + async embedQuery() { + return reference; + }, + async embedDocuments(texts) { + return texts.map(vector); + }, +}; +let server: ReturnType | undefined; +const percentile = (values: number[], fraction: number) => { + const sorted = [...values].sort((a, b) => a - b); + return Number( + sorted[ + Math.min(sorted.length - 1, Math.floor(sorted.length * fraction)) + ]!.toFixed(2), + ); +}; +try { + await admin.command({ query: `CREATE DATABASE ${database}` }); + await migrate(client, dimensions, true); + await store.probe(); + await client.command({ + query: + "CREATE TABLE retired (tenant_id String, namespace_id String, space_id String, file_id String, generation FixedString(32)) ENGINE=MergeTree ORDER BY (tenant_id,namespace_id,space_id,file_id,generation)", + }); + const start = performance.now(); + for (let offset = 0; offset < count; offset += 64) { + const rows: StoredChunk[] = Array.from( + { length: Math.min(64, count - offset) }, + (_, i) => ({ + ...scope, + fileId: "corpus", + generation, + spaceId, + index: offset + i, + text: `synthetic chunk ${offset + i}`, + page: null, + segment: 1, + start: 0, + end: 20, + section: [], + actor: "service:benchmark", + sourceClass: "asserted", + embedding: vector(), + }), + ); + await store.insert(rows, `bench-${offset}`, signal()); + } + const document: Document = { + ...scope, + fileId: "corpus", + generation, + spaceId, + version: "1", + state: "ready", + title: "Corpus", + original: null, + actor: "service:benchmark", + sourceClass: "asserted", + chunkCount: count, + operationKey: "a".repeat(64), + requestHash: "a".repeat(64), + }; + await store.publish(document, signal()); + const insertedMs = performance.now() - start; + const where = liveWhere([scope], spaceId); + const baselineWhere = + "tenant_id = {tenant:String} AND namespace_id = {lib0:String} AND space_id = {space:String} AND (tenant_id,namespace_id,space_id,file_id,generation) NOT IN (SELECT tenant_id,namespace_id,space_id,file_id,generation FROM retired)"; + const params = { ...where.params, vector: reference, k: 10 }; + const settings: ClickHouseSettings = { + vector_search_use_quantized_codes: "1", + vector_search_index_fetch_multiplier: "3", + use_query_condition_cache: 0, + use_query_cache: 0, + max_threads: 2, + }; + const times = { + baseline: [] as number[], + publication: [] as number[], + httpExact: [] as number[], + httpQuantized: [] as number[], + }; + for (let round = 0; round < 24; round++) { + for (const mode of (round % 2 + ? ["publication", "baseline"] + : ["baseline", "publication"]) as Array<"baseline" | "publication">) { + const started = performance.now(); + const result = await client.query({ + query: vectorSQL(mode === "baseline" ? baselineWhere : where.sql), + query_params: params, + format: "JSONEachRow", + clickhouse_settings: { + ...settings, + log_comment: `ragv2.benchmark.${mode}`, + }, + }); + await result.json(); + if (round >= 4) times[mode].push(performance.now() - started); + } + } + const exact = await client.query({ + query: vectorSQL(where.sql), + query_params: params, + format: "JSONEachRow", + clickhouse_settings: { + vector_search_use_quantized_codes: "0", + use_query_condition_cache: 0, + }, + }); + const nearest = await exact.json<{ chunk_index: number }>(); + const app = createApp({ auth, store, provider }); + server = Bun.serve({ hostname: "127.0.0.1", port: 0, fetch: app.fetch }); + const jwt = await token({ + tenant_id: "benchmark", + grants: [ + { + namespaceId: "library-a", + resourceKind: "document", + operations: ["read"], + }, + ], + }); + const recall = { exact: 0, quantized: 0 }; + for (let round = 0; round < 24; round++) { + for (const precision of ["exact", "quantized"] as const) { + const started = performance.now(); + const response = await fetch( + `http://127.0.0.1:${server.port}/v2/search`, + { + method: "POST", + headers: { + Authorization: `Bearer ${jwt}`, + "Content-Type": "application/json", + }, + body: JSON.stringify({ + query: "benchmark query", + namespaces: [{ namespaceId: "library-a" }], + k: 10, + precision, + }), + }, + ); + if (!response.ok) throw Error(`BENCHMARK_HTTP_${response.status}`); + const body = (await response.json()) as { + hits: Array<{ index: number }>; + }; + recall[precision] = + body.hits.filter((hit) => + nearest.some((exact) => exact.chunk_index === hit.index), + ).length / 10; + if (round >= 4) + times[precision === "exact" ? "httpExact" : "httpQuantized"].push( + performance.now() - started, + ); + } + } + console.log( + JSON.stringify({ + benchmark: "synthetic-single-namespace-warm-cache", + runtime: Bun.version, + dimensions, + rows: count, + insertMs: Number(insertedMs.toFixed(2)), + recallAt10: recall, + arms: Object.fromEntries( + Object.entries(times).map(([mode, values]) => [ + mode, + { p50Ms: percentile(values, 0.5), p95Ms: percentile(values, 0.95) }, + ]), + ), + note: "Baseline is Loom-shaped SQL over identical coded data, not a running Loom deployment. HTTP includes JWT verification, vector cache, candidate search, and metadata hydration. No external provider or extraction latency.", + }), + ); +} finally { + await server?.stop(true); + await store.close(); + await admin.command({ query: `DROP DATABASE IF EXISTS ${database} SYNC` }); + await admin.close(); +} diff --git a/v2/scripts/qualify.sh b/v2/scripts/qualify.sh new file mode 100755 index 00000000..c30307cd --- /dev/null +++ b/v2/scripts/qualify.sh @@ -0,0 +1,26 @@ +#!/usr/bin/env bash +set -euo pipefail +cd "$(dirname "$0")/.." +run_id="$(date +%s)-$$" +ch="ragv2-ch-$run_id" +runner="ragv2-bun-$run_id" +cleanup() { + docker rm -f "$runner" "$ch" >/dev/null 2>&1 || true +} +trap cleanup EXIT +# No published database port or application credentials. +docker run -d --name "$ch" --network none --cpus 2 --memory 2g \ + -e CLICKHOUSE_SKIP_USER_SETUP=1 clickhouse/clickhouse-server:26.8 >/dev/null +for attempt in $(seq 1 60); do + if docker exec "$ch" clickhouse-client --query 'SELECT 1' >/dev/null 2>&1; then break; fi + if [ "$attempt" = 60 ]; then printf 'ClickHouse did not become ready\n' >&2; exit 1; fi + sleep 1 +done +docker run -d --name "$runner" --network "container:$ch" --entrypoint sh \ + -e RAG_V2_TEST_CLICKHOUSE_URL=http://127.0.0.1:8123 oven/bun:1.4.2-debian -c 'sleep 300' >/dev/null +docker exec "$runner" mkdir /app +docker cp . "$runner:/app/v2" +docker exec -w /app/v2 "$runner" bun test test/integration.test.ts +if [ "${RAG_V2_RUN_BENCH:-false}" = true ]; then + docker exec -w /app/v2 "$runner" bun scripts/benchmark.ts +fi diff --git a/v2/src/app.ts b/v2/src/app.ts new file mode 100644 index 00000000..5e5a5ff6 --- /dev/null +++ b/v2/src/app.ts @@ -0,0 +1,297 @@ +import { Hono } from "hono"; +import { z } from "zod"; +import { randomUUID, createHash } from "node:crypto"; +import { authorize, type Authenticator, type Principal } from "./auth"; +import { + id, + ingestSchema, + originalSchema, + RagError, + searchSchema, + type Store, +} from "./contracts"; +import { QueryEmbeddings, type EmbeddingProvider } from "./embeddings"; +import type { Extractor } from "./extraction"; +import { Gate, readBounded } from "./limits"; +import { Pipeline } from "./pipeline"; +import { queryTag } from "./clickhouse"; + +type Environment = { Variables: { principal: Principal; signal: AbortSignal } }; +export function createApp(options: { + auth: Authenticator; + store: Store; + provider: EmbeddingProvider; + writer?: boolean; + extractor?: Extractor; +}) { + const app = new Hono(); + const requests = new Gate(64); + const searchRequests = new Gate(16); + const ingestionRequests = new Gate(2); + const queryEmbeddings = new QueryEmbeddings(options.provider); + const pipeline = new Pipeline(options.store, options.provider); + app.get("/health", (context) => context.json({ status: "ok", version: 2 })); + app.use("/v2/*", async (context, next) => { + if (context.req.header("X-API-Key") || context.req.header("User-Id")) + throw new RagError("AMBIGUOUS_AUTH", 400); + const signal = AbortSignal.any([ + context.req.raw.signal, + AbortSignal.timeout(30000), + ]); + return requests.run(signal, async () => { + const principal = await options.auth.verify( + context.req.header("Authorization"), + ); + context.set("principal", principal); + context.set("signal", signal); + const requestId = randomUUID(); + context.header("X-Request-Id", requestId); + context.header("Cache-Control", "no-store"); + await queryTag.run( + `${context.req.method} ${context.req.routePath} ${requestId}`, + next, + ); + }); + }); + const identity = (namespaceId: string, fileId: string) => { + id.parse(namespaceId); + id.parse(fileId); + }; + const writeKey = (value: string | undefined) => + z + .string() + .regex(/^[A-Za-z0-9._:-]{1,128}$/) + .parse(value); + const checkWriter = () => { + if (!options.writer) throw new RagError("WRITER_UNAVAILABLE", 503); + }; + app.put("/v2/namespaces/:namespaceId/documents/:fileId", async (context) => { + const { namespaceId, fileId } = context.req.param(); + identity(namespaceId, fileId); + const principal = context.get("principal"); + const scope = authorize(principal, namespaceId, "write", [fileId]); + const key = writeKey(context.req.header("Idempotency-Key")); + checkWriter(); + if ( + context.req.header("Content-Type")?.split(";")[0] !== "application/json" + ) + throw new RagError("INVALID_CONTENT_TYPE", 400); + return ingestionRequests.run(context.get("signal"), async () => { + const bytes = await readBounded( + context.req.raw.body, + 3 * 1024 * 1024, + context.get("signal"), + ); + let json: unknown; + try { + json = JSON.parse( + new TextDecoder("utf-8", { fatal: true }).decode(bytes), + ); + } catch { + throw new RagError("INVALID_BODY", 400); + } + const input = ingestSchema.parse(json); + if (input.original && input.original.fileId !== fileId) + throw new RagError("ORIGINAL_MISMATCH", 400); + const document = await pipeline.ingest( + scope, + fileId, + principal.sub, + key, + input, + context.get("signal"), + { actorKind: principal.actor_kind, sourceClass: "asserted" }, + ); + return context.json({ document, offsetEncoding: "utf16", pageBase: 1 }); + }); + }); + app.put( + "/v2/namespaces/:namespaceId/documents/:fileId/content", + async (context) => { + const { namespaceId, fileId } = context.req.param(); + identity(namespaceId, fileId); + const principal = context.get("principal"); + const scope = authorize(principal, namespaceId, "write", [fileId]); + const key = writeKey(context.req.header("Idempotency-Key")); + checkWriter(); + if (!options.extractor) throw new RagError("EXTRACTION_UNAVAILABLE", 503); + if (context.req.header("Content-Type") !== "application/octet-stream") + throw new RagError("INVALID_CONTENT_TYPE", 400); + const format = z + .enum(["pdf", "docx"]) + .parse(context.req.header("X-Rag-Format")); + const length = z + .string() + .regex(/^[1-9]\d{0,7}$/) + .parse(context.req.header("Content-Length")); + if (Number(length) > 10 * 1024 * 1024) + throw new RagError("BODY_LIMIT", 413); + const originalHeader = context.req.header("X-Rag-Original"); + if (!originalHeader || originalHeader.length > 2048) + throw new RagError("INVALID_ORIGINAL", 400); + let original: z.infer; + try { + original = originalSchema.parse(JSON.parse(originalHeader)); + } catch { + throw new RagError("INVALID_ORIGINAL", 400); + } + if (original.fileId !== fileId) + throw new RagError("ORIGINAL_MISMATCH", 400); + return ingestionRequests.run(context.get("signal"), async () => { + const bytes = await readBounded( + context.req.raw.body, + 10 * 1024 * 1024, + context.get("signal"), + ); + const digest = createHash("sha256").update(bytes).digest("hex"); + if (bytes.byteLength !== Number(length) || digest !== original.sha256) + throw new RagError("ORIGINAL_MISMATCH", 400); + const segments = await options.extractor!.extract( + bytes, + format, + digest, + context.get("signal"), + ); + const input = ingestSchema.parse({ segments, original }); + const document = await pipeline.ingest( + scope, + fileId, + principal.sub, + key, + input, + context.get("signal"), + { actorKind: principal.actor_kind, sourceClass: "extracted" }, + ); + return context.json({ document, offsetEncoding: "utf16", pageBase: 1 }); + }); + }, + ); + app.post("/v2/search", async (context) => + searchRequests.run(context.get("signal"), async () => { + if ( + context.req.header("Content-Type")?.split(";")[0] !== "application/json" + ) + throw new RagError("INVALID_CONTENT_TYPE", 400); + const bytes = await readBounded( + context.req.raw.body, + 65536, + context.get("signal"), + ); + let json: unknown; + try { + json = JSON.parse( + new TextDecoder("utf-8", { fatal: true }).decode(bytes), + ); + } catch { + throw new RagError("INVALID_BODY", 400); + } + const input = searchSchema.parse(json); + const scopes = input.namespaces.map((library) => + authorize( + context.get("principal"), + library.namespaceId, + "read", + library.resourceIds, + ), + ); + const vector = await queryEmbeddings.get( + input.query, + context.get("signal"), + ); + const hits = await options.store.search( + scopes, + vector, + options.provider.spaceId, + input, + context.get("signal"), + ); + return context.json({ + hits, + spaceId: options.provider.spaceId, + offsetEncoding: "utf16", + pageBase: 1, + }); + }), + ); + app.get("/v2/namespaces/:namespaceId/documents/:fileId", async (context) => { + const { namespaceId, fileId } = context.req.param(); + identity(namespaceId, fileId); + const scope = authorize(context.get("principal"), namespaceId, "read", [ + fileId, + ]); + const document = await options.store.get( + scope, + fileId, + context.get("signal"), + ); + if (!document || document.state !== "ready") + throw new RagError("NOT_FOUND", 404); + return context.json({ document }); + }); + app.get( + "/v2/namespaces/:namespaceId/documents/:fileId/context", + async (context) => { + const { namespaceId, fileId } = context.req.param(); + identity(namespaceId, fileId); + const scope = authorize(context.get("principal"), namespaceId, "read", [ + fileId, + ]); + const document = await options.store.get( + scope, + fileId, + context.get("signal"), + ); + if (!document || document.state !== "ready") + throw new RagError("NOT_FOUND", 404); + const chunks = await options.store.context( + scope, + document, + context.get("signal"), + ); + return context.json({ + document, + chunks, + offsetEncoding: "utf16", + pageBase: 1, + }); + }, + ); + app.delete( + "/v2/namespaces/:namespaceId/documents/:fileId", + async (context) => { + const { namespaceId, fileId } = context.req.param(); + identity(namespaceId, fileId); + const principal = context.get("principal"); + const scope = authorize(principal, namespaceId, "delete", [fileId]); + const key = writeKey(context.req.header("Idempotency-Key")); + checkWriter(); + await pipeline.delete( + scope, + fileId, + principal.sub, + key, + context.get("signal"), + principal.actor_kind, + ); + return context.body(null, 204); + }, + ); + app.onError((error, context) => { + if (error instanceof RagError) { + if (error.status === 503) context.header("Retry-After", "1"); + return context.json({ error: { code: error.code } }, error.status); + } + if (error instanceof z.ZodError) + return context.json({ error: { code: "INVALID_REQUEST" } }, 400); + if (context.get("signal")?.aborted) + return context.json({ error: { code: "DEADLINE" } }, 504); + console.error( + JSON.stringify({ + event: "ragv2.request_failed", + requestId: context.res.headers.get("X-Request-Id"), + }), + ); + return context.json({ error: { code: "BACKEND_UNAVAILABLE" } }, 503); + }); + return app; +} diff --git a/v2/src/auth.ts b/v2/src/auth.ts new file mode 100644 index 00000000..b2d011b2 --- /dev/null +++ b/v2/src/auth.ts @@ -0,0 +1,135 @@ +import { importJWK, jwtVerify, decodeProtectedHeader, type JWK } from "jose"; +import { z } from "zod"; +import { id, RagError, type Scope } from "./contracts"; + +const grantsSchema = z + .array( + z + .object({ + namespaceId: id, + resourceKind: z.literal("document"), + operations: z + .array(z.enum(["read", "write", "delete"])) + .min(1) + .max(3), + resourceIds: z.array(id).min(1).max(100).optional(), + }) + .strict(), + ) + .min(1) + .max(16); +const claimsSchema = z.object({ + sub: id, + actor_kind: z.enum(["user", "service"]), + tenant_id: id, + jti: z.string().min(1).max(256), + grants: grantsSchema, +}); +export type Principal = z.infer; +export interface Authenticator { + verify(authorization: string | undefined): Promise; +} + +export async function createAuthenticator( + jwks: { keys: JWK[] }, + issuer = "librechat", +): Promise { + if (!jwks.keys.length || jwks.keys.length > 16) + throw new Error("INVALID_VERIFICATION_KEYS"); + const keys = new Map< + string, + { alg: "EdDSA" | "RS256"; key: Awaited> } + >(); + for (const jwk of jwks.keys) { + if ( + !jwk.kid || + keys.has(jwk.kid) || + jwk.d || + jwk.k || + (jwk.alg !== "EdDSA" && jwk.alg !== "RS256") || + (jwk.alg === "EdDSA" && (jwk.kty !== "OKP" || jwk.crv !== "Ed25519")) || + (jwk.alg === "RS256" && jwk.kty !== "RSA") + ) + throw new Error("INVALID_VERIFICATION_KEYS"); + try { + keys.set(jwk.kid, { + alg: jwk.alg === "EdDSA" ? "EdDSA" : "RS256", + key: await importJWK(jwk, jwk.alg), + }); + } catch { + throw new Error("INVALID_VERIFICATION_KEYS"); + } + } + return { + async verify(authorization) { + if ( + !authorization || + authorization.length > 16384 || + !/^Bearer [\w.-]+$/.test(authorization) + ) + throw new RagError("AUTH_INVALID", 401); + try { + const token = authorization.slice(7); + const header = decodeProtectedHeader(token); + const entry = header.kid ? keys.get(header.kid) : undefined; + if (!entry || header.alg !== entry.alg || header.typ !== "JWT") + throw new Error("HEADER"); + const { payload } = await jwtVerify(token, entry.key, { + issuer, + audience: "ragapi", + algorithms: [entry.alg], + clockTolerance: 5, + maxTokenAge: 300, + requiredClaims: [ + "sub", + "actor_kind", + "tenant_id", + "jti", + "iat", + "nbf", + "exp", + "grants", + ], + }); + if ( + typeof payload.iat !== "number" || + typeof payload.exp !== "number" || + typeof payload.nbf !== "number" || + payload.exp <= payload.iat || + payload.exp - payload.iat > 300 || + payload.nbf > payload.exp + ) + throw new Error("LIFETIME"); + return claimsSchema.parse(payload); + } catch { + throw new RagError("AUTH_INVALID", 401); + } + }, + }; +} + +export function authorize( + principal: Principal, + namespaceId: string, + operation: "read" | "write" | "delete", + resourceIds?: readonly string[], +): Scope { + const grants = principal.grants.filter( + (grant) => + grant.namespaceId === namespaceId && grant.operations.includes(operation), + ); + if (!grants.length) throw new RagError("SCOPE_DENIED", 403); + const unrestricted = grants.some((grant) => grant.resourceIds === undefined); + const permitted = new Set(grants.flatMap((grant) => grant.resourceIds ?? [])); + if (resourceIds?.some((file) => !unrestricted && !permitted.has(file))) + throw new RagError("SCOPE_DENIED", 403); + return { + tenantId: principal.tenant_id, + namespaceId, + ...(!unrestricted || resourceIds + ? { + resourceIds: resourceIds ? [...new Set(resourceIds)] : [...permitted], + } + : {}), + }; +} diff --git a/v2/src/chunking.ts b/v2/src/chunking.ts new file mode 100644 index 00000000..f4c030b1 --- /dev/null +++ b/v2/src/chunking.ts @@ -0,0 +1,67 @@ +import type { Chunk, Segment } from "./contracts"; + +function boundary(text: string, position: number): number { + const code = text.charCodeAt(position - 1); + return code >= 0xd800 && code <= 0xdbff ? position - 1 : position; +} + +export function* chunks( + segments: readonly Segment[], + size = 1500, + overlap = 150, +): Generator { + if (size < 4 || overlap < 0 || overlap >= size - 2) + throw new Error("INVALID_CHUNK_SIZE"); + let index = 0; + for (const segment of segments) { + const sections: Array<{ start: number; end: number; path: string[] }> = []; + const stack: Array<{ depth: number; title: string }> = []; + let start = 0; + let path: string[] = []; + let offset = 0; + let fenced = false; + for (const line of segment.text.split("\n")) { + if (/^\s*(```|~~~)/.test(line)) fenced = !fenced; + const heading = fenced ? null : /^(#{1,6})[ \t]+(.+?)\s*$/.exec(line); + if (heading) { + sections.push({ start, end: offset, path }); + const depth = heading[1]!.length; + while (stack.length && stack.at(-1)!.depth >= depth) stack.pop(); + const title = heading[2]!; + stack.push({ + depth, + title: title.slice(0, boundary(title, Math.min(512, title.length))), + }); + path = stack.map((entry) => entry.title); + start = offset; + } + offset += line.length + 1; + } + sections.push({ start, end: segment.text.length, path }); + for (const section of sections) { + let position = section.start; + while (position < section.end) { + const end = boundary( + segment.text, + Math.min(position + size, section.end), + ); + const text = segment.text.slice(position, end); + if (text.trim()) + yield { + index: index++, + text, + page: segment.kind === "page" ? segment.index : null, + segment: segment.index, + start: position, + end, + section: section.path, + }; + if (end >= section.end) break; + position = Math.max( + position + 1, + boundary(segment.text, end - overlap), + ); + } + } + } +} diff --git a/v2/src/clickhouse.ts b/v2/src/clickhouse.ts new file mode 100644 index 00000000..7837bd96 --- /dev/null +++ b/v2/src/clickhouse.ts @@ -0,0 +1,572 @@ +import { + createClient, + TupleParam, + ClickHouseLogLevel, + type ClickHouseClient, + type ClickHouseSettings, +} from "@clickhouse/client"; +import { AsyncLocalStorage } from "node:async_hooks"; +import { createHash } from "node:crypto"; +import { Gate } from "./limits"; +import { + RagError, + narrowScope, + type Chunk, + type Document, + type Hit, + type Receipt, + type Scope, + type SearchInput, + type Store, + type StoredChunk, +} from "./contracts"; + +export const queryTag = new AsyncLocalStorage(); +export type DatabaseConfig = { + url: string; + username: string; + password: string; + database: string; +}; +export function databaseClient(config: DatabaseConfig): ClickHouseClient { + if (!/^[a-z][a-z0-9_]{0,62}$/.test(config.database)) + throw new Error("INVALID_DATABASE"); + return createClient({ + ...config, + application: "rag-api-v2-poc", + log: { level: ClickHouseLogLevel.OFF }, + max_open_connections: 12, + request_timeout: 15000, + keep_alive: { enabled: true, idle_socket_ttl: 2500 }, + compression: { request: true, response: true }, + clickhouse_settings: { + async_insert: 0, + wait_for_async_insert: 1, + output_format_json_quote_64bit_integers: 1, + max_execution_time: 10, + enable_filesystem_cache: 1, + }, + }); +} + +export function scopeWhere(scopes: readonly Scope[]): { + sql: string; + params: Record; +} { + if ( + !scopes.length || + scopes.some((scope) => scope.tenantId !== scopes[0]!.tenantId) + ) + throw new RagError("SCOPE_DENIED", 403); + const params: Record = { tenant: scopes[0]!.tenantId }; + const clauses = scopes.map((scope, index) => { + params[`lib${index}`] = scope.namespaceId; + if (scope.resourceIds !== undefined) + params[`files${index}`] = scope.resourceIds; + return `(namespace_id = {lib${index}:String}${scope.resourceIds !== undefined ? ` AND file_id IN {files${index}:Array(String)}` : ""})`; + }); + return { + sql: `tenant_id = {tenant:String} AND (${clauses.join(" OR ")})`, + params, + }; +} +export function liveWhere( + scopes: readonly Scope[], + spaceId: string, +): { sql: string; params: Record } { + const scope = scopeWhere(scopes); + return { + sql: `${scope.sql} AND space_id = {space:String} AND (tenant_id, namespace_id, file_id, generation) IN (SELECT tenant_id, namespace_id, file_id, generation FROM documents FINAL WHERE ${scope.sql} AND state = 'ready' AND space_id = {space:String})`, + params: { ...scope.params, space: spaceId }, + }; +} +export function vectorSQL(where: string): string { + return `SELECT namespace_id, file_id, generation, chunk_index, content, page, segment, char_start, char_end, section, + cosineDistance(embedding, {vector:Array(Float32)}) AS distance FROM chunks + WHERE ${where} ORDER BY distance ASC LIMIT {k:UInt32}`; +} +type DocumentRow = { + tenant_id: string; + namespace_id: string; + file_id: string; + generation: string; + version: string; + state: "ready" | "deleted"; + title: string; + original: string; + actor: string; + source_class: "asserted" | "extracted"; + space_id: string; + chunk_count: number; + operation_key: string; + request_hash: string; +}; +function hydrate(row: DocumentRow): Document { + return { + tenantId: row.tenant_id, + namespaceId: row.namespace_id, + fileId: row.file_id, + generation: row.generation, + version: row.version, + state: row.state, + title: row.title, + original: JSON.parse(row.original), + actor: row.actor, + sourceClass: row.source_class, + spaceId: row.space_id, + chunkCount: row.chunk_count, + operationKey: row.operation_key, + requestHash: row.request_hash, + }; +} +const documentColumns = + "tenant_id, namespace_id, file_id, generation, toString(version) AS version, state, title, original, actor, source_class, space_id, chunk_count, operation_key, request_hash"; +type HitRow = { + namespace_id: string; + file_id: string; + generation: string; + chunk_index: number; + content: string; + page: number | null; + segment: number; + char_start: number; + char_end: number; + section: string[]; + distance: number; + lexical?: number; +}; +const transient = (error: unknown) => + error instanceof Error && + /ECONNRESET|ECONNREFUSED|EPIPE|socket hang up/.test(error.message); + +export class ClickHouseStore implements Store { + private readonly reads = new Gate(8, 8); + private readonly writes = new Gate(2, 4); + private quantized = false; + private codesSettingSupported = false; + private probeAt = 0; + private probing: Promise | undefined; + get quantizedCodesEnabled(): boolean { + return this.quantized; + } + private readonly stats = new Map< + string, + { expires: number; n: number; avgdl: number; df: number[] } + >(); + constructor( + public readonly client: ClickHouseClient, + private readonly database: string, + public readonly dimensions: number, + ) {} + async probe(): Promise { + this.probing ??= this.probeState().finally(() => { + this.probing = undefined; + }); + return this.probing; + } + private async probeState(): Promise { + const signal = AbortSignal.timeout(5000); + try { + const rows = await this.rows<{ + supported: number; + coded: number; + patches: number; + }>( + `SELECT (SELECT count() FROM system.settings WHERE name = 'vector_search_use_quantized_codes') AS supported, + (SELECT count() FROM system.columns WHERE database = {db:String} AND table = 'chunks' AND name = 'embedding' AND compression_codec LIKE '%Quantized(%') AS coded, + (SELECT count() FROM system.parts WHERE database = {db:String} AND table = 'chunks' AND active AND startsWith(name, 'patch')) AS patches`, + { db: this.database }, + signal, + ); + this.codesSettingSupported = Number(rows[0]?.supported) > 0; + this.quantized = + Number(rows[0]?.supported) > 0 && + Number(rows[0]?.coded) > 0 && + Number(rows[0]?.patches) === 0; + if (this.quantized) { + try { + const replicas = await this.rows<{ + replicas: string; + supported: string; + }>( + "SELECT count() AS replicas, countIf(has_setting) AS supported FROM (SELECT max(name = 'vector_search_use_quantized_codes') AS has_setting FROM clusterAllReplicas('default', system.settings) GROUP BY hostName())", + {}, + signal, + ); + this.quantized = + Number(replicas[0]?.replicas) > 0 && + replicas[0]?.supported === replicas[0]?.replicas; + } catch { + /* Standalone OSS has no default cluster. Reads still handle unknown settings. */ + } + } + } catch { + this.quantized = false; + this.codesSettingSupported = false; + } + this.probeAt = Date.now(); + } + private settings(extra: ClickHouseSettings = {}): ClickHouseSettings { + return { + log_comment: (queryTag.getStore() ?? "ragv2.background").slice(0, 120), + ...extra, + }; + } + async rows( + query: string, + params: Record, + signal: AbortSignal, + extra: ClickHouseSettings = {}, + ): Promise { + return this.reads.run(signal, async () => { + for (let attempt = 0; ; attempt++) { + try { + const result = await this.client.query({ + query, + query_params: params, + format: "JSONEachRow", + abort_signal: signal, + clickhouse_settings: this.settings(extra), + }); + try { + return await result.json(); + } finally { + result.close(); + } + } catch (error) { + signal.throwIfAborted(); + if (attempt || !transient(error)) throw error; + } + } + }); + } + private async write( + table: string, + values: readonly Record[], + token: string, + signal: AbortSignal, + ): Promise { + await this.writes.run(signal, async () => { + for (let attempt = 0; ; attempt++) { + try { + await this.client.insert({ + table, + values, + format: "JSONEachRow", + abort_signal: signal, + clickhouse_settings: this.settings({ + async_insert: 0, + wait_for_async_insert: 1, + insert_deduplicate: 1, + insert_deduplication_token: token, + }), + }); + return; + } catch (error) { + signal.throwIfAborted(); + if (attempt || !transient(error)) throw error; + } + } + }); + } + async get( + scope: Scope, + fileId: string, + signal: AbortSignal, + ): Promise { + const where = scopeWhere([narrowScope(scope, fileId)]); + const rows = await this.rows( + `SELECT ${documentColumns} FROM documents FINAL WHERE ${where.sql} LIMIT 1`, + where.params, + signal, + { select_sequential_consistency: "1" }, + ); + return rows[0] ? hydrate(rows[0]) : null; + } + async receipt( + scope: Scope, + operationKey: string, + signal: AbortSignal, + ): Promise { + const rows = await this.rows<{ receipt: string }>( + "SELECT receipt FROM operations FINAL WHERE tenant_id = {tenant:String} AND namespace_id = {library:String} AND operation_key = {key:String} LIMIT 1", + { tenant: scope.tenantId, library: scope.namespaceId, key: operationKey }, + signal, + { select_sequential_consistency: "1" }, + ); + return rows[0] ? JSON.parse(rows[0].receipt) : null; + } + async insert( + chunks: readonly StoredChunk[], + batchId: string, + signal: AbortSignal, + ): Promise { + await this.write( + "chunks", + chunks.map((chunk) => ({ + tenant_id: chunk.tenantId, + namespace_id: chunk.namespaceId, + file_id: chunk.fileId, + generation: chunk.generation, + space_id: chunk.spaceId, + chunk_index: chunk.index, + content: chunk.text, + page: chunk.page, + segment: chunk.segment, + char_start: chunk.start, + char_end: chunk.end, + section: chunk.section, + actor: chunk.actor, + source_class: chunk.sourceClass, + embedding: chunk.embedding, + })), + batchId, + signal, + ); + } + async publish(document: Document, signal: AbortSignal): Promise { + await this.write( + "documents", + [ + { + tenant_id: document.tenantId, + namespace_id: document.namespaceId, + file_id: document.fileId, + generation: document.generation, + version: document.version, + state: document.state, + title: document.title, + original: JSON.stringify(document.original), + actor: document.actor, + source_class: document.sourceClass, + space_id: document.spaceId, + chunk_count: document.chunkCount, + operation_key: document.operationKey, + request_hash: document.requestHash, + }, + ], + `${document.operationKey}:${document.version}`, + signal, + ); + this.stats.clear(); + } + async record( + scope: Scope, + operationKey: string, + receipt: Receipt, + signal: AbortSignal, + ): Promise { + await this.write( + "operations", + [ + { + tenant_id: scope.tenantId, + namespace_id: scope.namespaceId, + operation_key: operationKey, + receipt: JSON.stringify(receipt), + version: receipt.document.version, + }, + ], + `${operationKey}:${receipt.document.version}`, + signal, + ); + } + async search( + scopes: readonly Scope[], + vector: readonly number[], + spaceId: string, + input: SearchInput, + signal: AbortSignal, + ): Promise { + if (Date.now() - this.probeAt > 60000) await this.probe(); + const where = liveWhere(scopes, spaceId); + const tokens = [ + ...new Set(input.query.toLowerCase().match(/[a-z0-9]{2,}/g) ?? []), + ].slice(0, 16); + const k = input.mode === "hybrid" ? Math.max(input.k * 3, 20) : input.k; + const params = { ...where.params, vector, k }; + const vectorQuery = async () => { + try { + return await this.rows(vectorSQL(where.sql), params, signal, { + select_sequential_consistency: "1", + use_query_condition_cache: 0, + ...(this.quantized && input.precision === "quantized" + ? { + vector_search_use_quantized_codes: "1", + vector_search_index_fetch_multiplier: "3", + } + : this.codesSettingSupported + ? { vector_search_use_quantized_codes: "0" } + : {}), + }); + } catch (error) { + if ( + !this.codesSettingSupported || + !(error instanceof Error) || + !( + ("code" in error && String(error.code) === "115") || + /Unknown setting|UNKNOWN_SETTING/i.test(error.message) + ) + ) + throw error; + this.quantized = false; + this.codesSettingSupported = false; + return this.rows(vectorSQL(where.sql), params, signal, { + select_sequential_consistency: "1", + use_query_condition_cache: 0, + }); + } + }; + const [semantic, lexical] = await Promise.all([ + vectorQuery(), + input.mode === "hybrid" && tokens.length + ? this.lexical(where, params, tokens, signal) + : Promise.resolve([]), + ]); + const byKey = new Map(); + const keyOf = (row: HitRow) => + JSON.stringify([ + row.namespace_id, + row.file_id, + row.generation, + row.chunk_index, + ]); + for (const [index, row] of semantic.entries()) + byKey.set(keyOf(row), { + row, + score: lexical.length ? 1 / (60 + index + 1) : 1 - row.distance, + }); + for (const [index, row] of lexical.entries()) { + const key = keyOf(row); + const prior = byKey.get(key); + byKey.set(key, { + row: prior?.row ?? row, + score: (prior?.score ?? 0) + 1 / (60 + index + 1), + }); + } + const top = [...byKey.values()] + .sort((a, b) => b.score - a.score || a.row.distance - b.row.distance) + .slice(0, input.k); + if (!top.length) return []; + const scoped = scopeWhere(scopes); + const docs = await this.rows( + `SELECT ${documentColumns} FROM documents FINAL WHERE ${scoped.sql} AND state = 'ready' AND (namespace_id, file_id, generation) IN {keys:Array(Tuple(String, String, String))}`, + { + ...scoped.params, + keys: top.map( + ({ row }) => + new TupleParam([row.namespace_id, row.file_id, row.generation]), + ), + }, + signal, + { select_sequential_consistency: "1" }, + ); + const metadata = new Map( + docs.map((row) => [ + JSON.stringify([row.namespace_id, row.file_id, row.generation]), + hydrate(row), + ]), + ); + return top.flatMap(({ row, score }) => { + const doc = metadata.get( + JSON.stringify([row.namespace_id, row.file_id, row.generation]), + ); + return doc + ? [ + { + namespaceId: row.namespace_id, + fileId: row.file_id, + generation: row.generation, + index: row.chunk_index, + text: row.content, + page: row.page, + segment: row.segment, + start: row.char_start, + end: row.char_end, + section: row.section, + distance: row.distance, + score, + original: doc.original, + actor: doc.actor, + sourceClass: doc.sourceClass, + }, + ] + : []; + }); + } + private async lexical( + where: ReturnType, + params: Record, + tokens: string[], + signal: AbortSignal, + ): Promise { + const key = createHash("sha256") + .update(JSON.stringify([where.params, tokens])) + .digest("hex"); + let stats = this.stats.get(key); + if (!stats || stats.expires <= Date.now()) { + const tokenParams = Object.fromEntries( + tokens.map((token, index) => [`term${index}`, token]), + ); + const dfSQL = tokens + .map((_, index) => `countIf(has(words, {term${index}:String}))`) + .join(", "); + const rows = await this.rows<{ n: string; avgdl: number; df: string[] }>( + `WITH splitByNonAlpha(lowerUTF8(content)) AS words SELECT count() AS n, avg(length(words)) AS avgdl, [${dfSQL}] AS df FROM chunks WHERE ${where.sql}`, + { ...params, ...tokenParams }, + signal, + { select_sequential_consistency: "1" }, + ); + const row = rows[0]; + if (!row || !Number(row.n) || !row.avgdl) return []; + stats = { + n: Number(row.n), + avgdl: row.avgdl, + df: row.df.map(Number), + expires: Date.now() + 30000, + }; + this.stats.set(key, stats); + while (this.stats.size > 128) + this.stats.delete(this.stats.keys().next().value!); + } + const idfs = stats.df.map((df) => + Math.log(1 + (stats.n - df + 0.5) / (df + 0.5)), + ); + return this.rows( + `WITH splitByNonAlpha(content_lc) AS words SELECT namespace_id, file_id, generation, chunk_index, content, page, segment, char_start, char_end, section, cosineDistance(embedding, {vector:Array(Float32)}) AS distance, + arraySum(arrayMap((tok, idf) -> idf * countEqual(words, tok) * 2.2 / (countEqual(words, tok) + 1.2 * (0.25 + 0.75 * length(words) / {avgdl:Float32})), {tokens:Array(String)}, {idfs:Array(Float32)})) AS lexical + FROM chunks WHERE ${where.sql} AND hasAnyTokens(content_lc, {tokens:Array(String)}) ORDER BY lexical DESC, distance ASC LIMIT {k:UInt32}`, + { ...params, tokens, idfs, avgdl: stats.avgdl }, + signal, + { select_sequential_consistency: "1" }, + ); + } + async context( + scope: Scope, + document: Document, + signal: AbortSignal, + ): Promise { + const where = liveWhere( + [narrowScope(scope, document.fileId)], + document.spaceId, + ); + const rows = await this.rows( + `SELECT chunk_index, content, page, segment, char_start, char_end, section FROM chunks WHERE ${where.sql} AND generation = {generation:String} ORDER BY chunk_index ASC LIMIT 10000`, + { ...where.params, generation: document.generation }, + signal, + { select_sequential_consistency: "1" }, + ); + return rows.map((row) => ({ + index: row.chunk_index, + text: row.content, + page: row.page, + segment: row.segment, + start: row.char_start, + end: row.char_end, + section: row.section, + })); + } + close(): Promise { + return this.client.close(); + } +} diff --git a/v2/src/config.ts b/v2/src/config.ts new file mode 100644 index 00000000..81e5d458 --- /dev/null +++ b/v2/src/config.ts @@ -0,0 +1,59 @@ +import { z } from "zod"; +import { + deterministicProvider, + openAICompatible, + type EmbeddingProvider, +} from "./embeddings"; +import type { DatabaseConfig } from "./clickhouse"; + +export function databaseConfig(): DatabaseConfig { + const url = process.env.RAG_V2_CLICKHOUSE_URL ?? "http://localhost:8123"; + const parsed = new URL(url); + if ( + parsed.username || + parsed.password || + parsed.search || + parsed.hash || + (parsed.protocol !== "https:" && + !( + parsed.protocol === "http:" && + (["localhost", "127.0.0.1", "[::1]"].includes(parsed.hostname) || + process.env.RAG_V2_ALLOW_INSECURE_CLICKHOUSE === "true") + )) + ) + throw new Error("INVALID_CLICKHOUSE_URL"); + return { + url, + username: process.env.RAG_V2_CLICKHOUSE_USER ?? "default", + password: process.env.RAG_V2_CLICKHOUSE_PASSWORD ?? "", + database: process.env.RAG_V2_CLICKHOUSE_DATABASE ?? "rag_v2", + }; +} +export function embeddingProvider(): EmbeddingProvider { + const dimensions = z.coerce + .number() + .int() + .min(8) + .max(8192) + .parse(process.env.RAG_V2_EMBEDDING_DIMENSIONS ?? 1536); + if (process.env.RAG_V2_EMBEDDING_PROVIDER === "test") { + if (process.env.RAG_V2_ALLOW_TEST_PROVIDER !== "true") + throw new Error("TEST_PROVIDER_NOT_ALLOWED"); + return deterministicProvider(dimensions); + } + if ( + process.env.RAG_V2_EMBEDDING_PROVIDER && + process.env.RAG_V2_EMBEDDING_PROVIDER !== "openai-compatible" + ) + throw new Error("INVALID_EMBEDDING_PROVIDER"); + const apiKey = process.env.RAG_V2_EMBEDDING_API_KEY; + if (!apiKey) throw new Error("EMBEDDING_KEY_REQUIRED"); + return openAICompatible({ + endpoint: + process.env.RAG_V2_EMBEDDING_URL ?? + "https://api.openai.com/v1/embeddings", + apiKey, + model: process.env.RAG_V2_EMBEDDING_MODEL ?? "text-embedding-3-small", + dimensions, + }); +} diff --git a/v2/src/contracts.ts b/v2/src/contracts.ts new file mode 100644 index 00000000..ce1c7d23 --- /dev/null +++ b/v2/src/contracts.ts @@ -0,0 +1,194 @@ +import { z } from "zod"; + +export const id = z.string().regex(/^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/); +export const originalSchema = z + .object({ + fileId: id, + revision: id, + sha256: z.string().regex(/^[a-f0-9]{64}$/), + filename: z.string().min(1).max(255), + mediaType: z + .string() + .min(1) + .max(128) + .regex(/^[\w.+-]+\/[\w.+-]+$/), + }) + .strict(); +export const segmentSchema = z + .object({ + kind: z.enum(["page", "document"]), + index: z.number().int().min(1).max(128), + text: z + .string() + .max(1024 * 1024) + .refine((text) => text.isWellFormed()), + }) + .strict(); +export const ingestSchema = z + .object({ + title: z.string().max(512).default(""), + segments: z.array(segmentSchema).min(1).max(128), + original: originalSchema.nullable().default(null), + ifMatch: z + .string() + .regex(/^[a-f0-9]{32}$/) + .optional(), + }) + .strict() + .superRefine((input, context) => { + let bytes = 0; + const first = input.segments[0]!; + for (const [index, segment] of input.segments.entries()) { + bytes += Buffer.byteLength(segment.text); + if (segment.kind !== first.kind || segment.index !== index + 1) { + context.addIssue({ + code: "custom", + message: "Segments must be ordered and homogeneous", + }); + } + } + if (first.kind === "document" && input.segments.length !== 1) { + context.addIssue({ + code: "custom", + message: "A document has one segment", + }); + } + if ( + bytes > 1024 * 1024 || + !input.segments.some((segment) => segment.text.trim()) + ) { + context.addIssue({ + code: "custom", + message: "Text must be nonempty and at most 1 MiB", + }); + } + }); +export const searchSchema = z + .object({ + query: z.string().trim().min(1).max(8192), + namespaces: z + .array( + z + .object({ + namespaceId: id, + resourceIds: z.array(id).min(1).max(100).optional(), + }) + .strict(), + ) + .min(1) + .max(8), + k: z.number().int().min(1).max(100).default(5), + mode: z.enum(["semantic", "hybrid"]).default("semantic"), + precision: z.enum(["exact", "quantized"]).default("exact"), + }) + .strict(); +export type IngestInput = z.infer; +export type Original = z.infer; +export type Segment = z.infer; +export type SearchInput = z.infer; +export type Scope = { + tenantId: string; + namespaceId: string; + resourceIds?: readonly string[]; +}; +export type Document = { + tenantId: string; + namespaceId: string; + fileId: string; + generation: string; + version: string; + state: "ready" | "deleted"; + title: string; + original: Original | null; + actor: string; + sourceClass: "asserted" | "extracted"; + spaceId: string; + chunkCount: number; + operationKey: string; + requestHash: string; +}; +export type Chunk = { + index: number; + text: string; + page: number | null; + segment: number; + start: number; + end: number; + section: string[]; +}; +export type StoredChunk = Chunk & { + tenantId: string; + namespaceId: string; + fileId: string; + generation: string; + spaceId: string; + actor: string; + sourceClass: "asserted" | "extracted"; + embedding: readonly number[]; +}; +export type Hit = Chunk & { + namespaceId: string; + fileId: string; + generation: string; + distance: number; + score: number; + original: Original | null; + actor: string; + sourceClass: "asserted" | "extracted"; +}; +export type Receipt = { requestHash: string; document: Document }; +export interface Store { + get( + scope: Scope, + fileId: string, + signal: AbortSignal, + ): Promise; + receipt( + scope: Scope, + operationKey: string, + signal: AbortSignal, + ): Promise; + insert( + chunks: readonly StoredChunk[], + batchId: string, + signal: AbortSignal, + ): Promise; + publish(document: Document, signal: AbortSignal): Promise; + record( + scope: Scope, + operationKey: string, + receipt: Receipt, + signal: AbortSignal, + ): Promise; + search( + scopes: readonly Scope[], + vector: readonly number[], + spaceId: string, + input: SearchInput, + signal: AbortSignal, + ): Promise; + context( + scope: Scope, + document: Document, + signal: AbortSignal, + ): Promise; + close(): Promise; +} +export class RagError extends Error { + constructor( + public readonly code: string, + public readonly status: + 400 | 401 | 403 | 404 | 409 | 413 | 422 | 429 | 500 | 503 | 504, + ) { + super(code); + } +} + +export function narrowScope(scope: Scope, resourceId: string): Scope { + if ( + scope.resourceIds !== undefined && + !scope.resourceIds.includes(resourceId) + ) + throw new RagError("SCOPE_DENIED", 403); + return { ...scope, resourceIds: [resourceId] }; +} diff --git a/v2/src/embeddings.ts b/v2/src/embeddings.ts new file mode 100644 index 00000000..f6f2c53b --- /dev/null +++ b/v2/src/embeddings.ts @@ -0,0 +1,223 @@ +import { createHash } from "node:crypto"; +import { z } from "zod"; +import { Gate, readBounded } from "./limits"; +import { RagError } from "./contracts"; + +export interface EmbeddingProvider { + readonly spaceId: string; + readonly dimensions: number; + embedQuery(text: string, signal: AbortSignal): Promise; + embedDocuments( + texts: readonly string[], + signal: AbortSignal, + ): Promise; +} +export function validateVectors( + vectors: readonly (readonly number[])[], + count: number, + dimensions: number, +): void { + if ( + vectors.length !== count || + vectors.some( + (vector) => + vector.length !== dimensions || + vector.some((value) => !Number.isFinite(value)) || + !vector.some((value) => value !== 0), + ) + ) + throw new RagError("INVALID_EMBEDDINGS", 503); +} +const responseSchema = z.object({ + data: z.array( + z.object({ + index: z.number().int().nonnegative(), + embedding: z.array(z.number().finite()), + }), + ), +}); + +export function openAICompatible(options: { + endpoint: string; + apiKey: string; + model: string; + dimensions: number; + concurrency?: number; +}): EmbeddingProvider { + const url = new URL(options.endpoint); + if ( + url.protocol !== "https:" && + !( + url.protocol === "http:" && + ["localhost", "127.0.0.1", "[::1]"].includes(url.hostname) + ) + ) + throw new Error("EMBEDDING_ENDPOINT_REQUIRES_TLS"); + if (url.username || url.password || url.hash || url.search) + throw new Error("INVALID_EMBEDDING_ENDPOINT"); + const spaceId = createHash("sha256") + .update( + JSON.stringify([ + url.href, + options.model, + options.dimensions, + "raw-query/title-section-document-v1", + ]), + ) + .digest("hex"); + const gate = new Gate(options.concurrency ?? 4, 16); + const embed = (texts: readonly string[], signal: AbortSignal) => + gate.run(signal, async () => { + if ( + !texts.length || + texts.length > 64 || + texts.some((text) => !text.trim() || Buffer.byteLength(text) > 8191) + ) + throw new RagError("EMBEDDING_INPUT_LIMIT", 422); + for (let attempt = 0; attempt < 2; attempt++) { + let response: Response; + try { + response = await fetch(url, { + method: "POST", + redirect: "error", + signal, + headers: { + Authorization: `Bearer ${options.apiKey}`, + "Content-Type": "application/json", + }, + body: JSON.stringify({ + model: options.model, + dimensions: options.dimensions, + input: texts, + }), + }); + } catch (error) { + signal.throwIfAborted(); + if (attempt || !(error instanceof TypeError)) + throw new RagError("EMBEDDING_UNAVAILABLE", 503); + await Bun.sleep(100); + continue; + } + if (response.status === 429 || response.status >= 500) { + await response.body?.cancel(); + if (attempt) throw new RagError("EMBEDDING_UNAVAILABLE", 503); + await Bun.sleep(100); + signal.throwIfAborted(); + continue; + } + if (!response.ok) { + await response.body?.cancel(); + throw new RagError("EMBEDDING_REJECTED", 503); + } + let parsed: z.infer; + try { + parsed = responseSchema.parse( + JSON.parse( + new TextDecoder("utf-8", { fatal: true }).decode( + await readBounded( + response.body, + texts.length * options.dimensions * 32 + 65536, + signal, + ), + ), + ), + ); + } catch { + signal.throwIfAborted(); + throw new RagError("INVALID_EMBEDDINGS", 503); + } + const vectors: number[][] = Array(texts.length); + for (const entry of parsed.data) { + if (entry.index >= texts.length || vectors[entry.index]) + throw new RagError("INVALID_EMBEDDINGS", 503); + vectors[entry.index] = entry.embedding; + } + if ( + parsed.data.length !== texts.length || + vectors.some((vector) => !vector) + ) + throw new RagError("INVALID_EMBEDDINGS", 503); + validateVectors(vectors, texts.length, options.dimensions); + return vectors; + } + throw new RagError("EMBEDDING_UNAVAILABLE", 503); + }); + return { + spaceId, + dimensions: options.dimensions, + embedQuery: async (text, signal) => (await embed([text], signal))[0]!, + embedDocuments: embed, + }; +} + +export function deterministicProvider(dimensions = 64): EmbeddingProvider { + const embed = (text: string) => { + const vector = Array(dimensions).fill(0); + for (const token of text.toLowerCase().match(/[a-z0-9]+/g) ?? []) { + const hash = createHash("sha256").update(token).digest(); + const index = hash.readUInt32LE(0) % dimensions; + vector[index] = vector[index]! + 1; + } + if (!vector.some(Boolean)) vector[0] = 1; + return vector; + }; + return { + spaceId: `test-only-token-hash-${dimensions}-v1`, + dimensions, + async embedQuery(text, signal) { + signal.throwIfAborted(); + return embed(text); + }, + async embedDocuments(texts, signal) { + signal.throwIfAborted(); + return texts.map(embed); + }, + }; +} + +export class QueryEmbeddings { + private readonly cache = new Map< + string, + { expires: number; vector: readonly number[] } + >(); + private readonly pending = new Map>(); + constructor( + private readonly provider: EmbeddingProvider, + private readonly maxEntries = 128, + ) {} + async get(text: string, signal: AbortSignal): Promise { + signal.throwIfAborted(); + const key = createHash("sha256") + .update(`${this.provider.spaceId}\0${text}`) + .digest("hex"); + const cached = this.cache.get(key); + if (cached && cached.expires > Date.now()) { + this.cache.delete(key); + this.cache.set(key, cached); + return cached.vector; + } + this.cache.delete(key); + let task = this.pending.get(key); + if (!task) { + if (this.pending.size >= 16) throw new RagError("OVERLOADED", 503); + task = this.provider + .embedQuery(text, AbortSignal.timeout(10000)) + .then((vector) => { + validateVectors([vector], 1, this.provider.dimensions); + this.cache.set(key, { vector, expires: Date.now() + 60000 }); + while (this.cache.size > this.maxEntries) + this.cache.delete(this.cache.keys().next().value!); + return vector; + }) + .finally(() => this.pending.delete(key)); + this.pending.set(key, task); + } + return new Promise((resolve, reject) => { + const abort = () => reject(signal.reason); + signal.addEventListener("abort", abort, { once: true }); + task + .then(resolve, reject) + .finally(() => signal.removeEventListener("abort", abort)); + }); + } +} diff --git a/v2/src/extraction.ts b/v2/src/extraction.ts new file mode 100644 index 00000000..adc9835d --- /dev/null +++ b/v2/src/extraction.ts @@ -0,0 +1,85 @@ +import { z } from "zod"; +import { ingestSchema, RagError, type Segment } from "./contracts"; +import { Gate, readBounded } from "./limits"; + +export interface Extractor { + extract( + bytes: Uint8Array, + format: "pdf" | "docx", + digest: string, + signal: AbortSignal, + ): Promise; +} +const resultSchema = z + .object({ + version: z.literal(1), + operation: z.literal("document.extract-text"), + format: z.enum(["pdf", "docx"]), + inputSha256: z.string(), + textBytes: z.number().int().nonnegative(), + segments: ingestSchema.shape.segments, + }) + .strict(); + +export function unixExtractor(socket: string): Extractor { + if (!socket.startsWith("/") || socket.includes("\0")) + throw new Error("INVALID_EXTRACTION_SOCKET"); + const gate = new Gate(2); + return { + extract: (bytes, format, digest, signal) => + gate.run(signal, async () => { + let response: Response; + try { + response = await fetch("http://localhost/v1/extract-text", { + unix: socket, + method: "POST", + redirect: "error", + signal: AbortSignal.any([signal, AbortSignal.timeout(11000)]), + headers: { + "Content-Type": "application/octet-stream", + "Content-Length": String(bytes.byteLength), + "X-Extraction-Version": "1", + "X-Extraction-Format": format, + }, + body: bytes, + }); + } catch { + signal.throwIfAborted(); + throw new RagError("EXTRACTION_UNAVAILABLE", 503); + } + const wire = await readBounded(response.body, 3 * 1024 * 1024, signal); + if (!response.ok) { + if (response.status === 429 || response.status === 503) + throw new RagError("EXTRACTION_BUSY", 503); + if (response.status === 504) + throw new RagError("EXTRACTION_DEADLINE", 504); + throw new RagError( + "EXTRACTION_REJECTED", + response.status === 413 ? 413 : 422, + ); + } + try { + const result = resultSchema.parse( + JSON.parse(new TextDecoder("utf-8", { fatal: true }).decode(wire)), + ); + if ( + result.format !== format || + result.inputSha256 !== digest || + result.textBytes !== + result.segments.reduce( + (size, segment) => size + Buffer.byteLength(segment.text), + 0, + ) || + result.segments.some( + (segment) => + segment.kind !== (format === "pdf" ? "page" : "document"), + ) + ) + throw new Error("INVALID_RESULT"); + return ingestSchema.parse({ segments: result.segments }).segments; + } catch { + throw new RagError("INVALID_EXTRACTION_RESULT", 503); + } + }), + }; +} diff --git a/v2/src/limits.ts b/v2/src/limits.ts new file mode 100644 index 00000000..c0384a3a --- /dev/null +++ b/v2/src/limits.ts @@ -0,0 +1,85 @@ +import { RagError } from "./contracts"; + +export class Gate { + private active = 0; + private readonly waiting: Array<() => void> = []; + constructor( + private readonly capacity: number, + private readonly queueSize = 0, + ) { + if (!Number.isInteger(capacity) || capacity < 1) + throw new Error("INVALID_CAPACITY"); + } + async run(signal: AbortSignal, work: () => Promise): Promise { + signal.throwIfAborted(); + if (this.active >= this.capacity) { + if (this.waiting.length >= this.queueSize) + throw new RagError("OVERLOADED", 503); + await new Promise((resolve, reject) => { + const take = () => { + signal.removeEventListener("abort", cancel); + resolve(); + }; + const cancel = () => { + const index = this.waiting.indexOf(take); + if (index >= 0) this.waiting.splice(index, 1); + reject(signal.reason); + }; + this.waiting.push(take); + signal.addEventListener("abort", cancel, { once: true }); + }); + } else this.active++; + try { + signal.throwIfAborted(); + return await work(); + } finally { + const next = this.waiting.shift(); + if (next) next(); + else this.active--; + } + } +} + +export class DocumentLocks { + private readonly active = new Set(); + async run(key: string, work: () => Promise): Promise { + if (this.active.has(key)) throw new RagError("DOCUMENT_BUSY", 409); + this.active.add(key); + try { + return await work(); + } finally { + this.active.delete(key); + } + } +} + +export async function readBounded( + body: ReadableStream | null, + limit: number, + signal: AbortSignal, +): Promise> { + if (!body) throw new RagError("INVALID_BODY", 400); + const reader = body.getReader(); + const parts: Uint8Array[] = []; + let size = 0; + const cancel = () => { + void reader.cancel().catch(() => {}); + }; + signal.addEventListener("abort", cancel, { once: true }); + try { + while (true) { + signal.throwIfAborted(); + const { value, done } = await reader.read(); + signal.throwIfAborted(); + if (done) break; + size += value.byteLength; + if (size > limit) throw new RagError("BODY_LIMIT", 413); + parts.push(value); + } + return Buffer.concat(parts, size); + } finally { + signal.removeEventListener("abort", cancel); + await reader.cancel().catch(() => {}); + reader.releaseLock(); + } +} diff --git a/v2/src/main.ts b/v2/src/main.ts new file mode 100644 index 00000000..df9a5838 --- /dev/null +++ b/v2/src/main.ts @@ -0,0 +1,77 @@ +import { z } from "zod"; +import { createApp } from "./app"; +import { createAuthenticator } from "./auth"; +import { ClickHouseStore, databaseClient } from "./clickhouse"; +import { databaseConfig, embeddingProvider } from "./config"; +import { unixExtractor } from "./extraction"; + +try { + const jwks = JSON.parse(process.env.RAG_V2_JWKS_JSON ?? "null"); + if (!jwks || !Array.isArray(jwks.keys)) + throw new Error("AUTH_CONFIGURATION_REQUIRED"); + const auth = await createAuthenticator( + jwks, + process.env.RAG_V2_JWT_ISSUER ?? "librechat", + ); + const provider = embeddingProvider(); + const config = databaseConfig(); + const store = new ClickHouseStore( + databaseClient(config), + config.database, + provider.dimensions, + ); + await store.rows("SELECT 1", {}, AbortSignal.timeout(5000)); + await store.probe(); + const mode = z + .enum(["api", "coordinator"]) + .parse(process.env.RAG_V2_MODE ?? "api"); + // Exactly one coordinator per database. API replicas are read-only. + if (mode === "coordinator" && process.env.RAG_V2_SINGLE_WRITER !== "true") + throw new Error("SINGLE_WRITER_ACKNOWLEDGMENT_REQUIRED"); + const extractor = process.env.RAG_V2_EXTRACTION_SOCKET + ? unixExtractor(process.env.RAG_V2_EXTRACTION_SOCKET) + : undefined; + const app = createApp({ + auth, + provider, + store, + writer: mode === "coordinator", + extractor, + }); + const server = Bun.serve({ + hostname: process.env.RAG_V2_HOST ?? "127.0.0.1", + port: z.coerce + .number() + .int() + .min(1) + .max(65535) + .parse(process.env.RAG_V2_PORT ?? 8001), + fetch: app.fetch, + maxRequestBodySize: 10 * 1024 * 1024, + }); + console.log( + JSON.stringify({ + event: "ragv2.ready", + mode, + port: server.port, + spaceId: provider.spaceId, + }), + ); + let stopping = false; + const shutdown = async () => { + if (stopping) return; + stopping = true; + const deadline = setTimeout(() => { + void server.stop(true); + process.exit(1); + }, 35000); + await server.stop(false); + await store.close(); + clearTimeout(deadline); + }; + process.once("SIGTERM", () => void shutdown()); + process.once("SIGINT", () => void shutdown()); +} catch { + console.error("RAG_V2_STARTUP_FAILED"); + process.exit(1); +} diff --git a/v2/src/migrate.ts b/v2/src/migrate.ts new file mode 100644 index 00000000..4404c56b --- /dev/null +++ b/v2/src/migrate.ts @@ -0,0 +1,32 @@ +import { databaseClient } from "./clickhouse"; +import { databaseConfig, embeddingProvider } from "./config"; +import { migrate } from "./schema"; + +const config = databaseConfig(); +const client = databaseClient(config); +try { + const provider = embeddingProvider(); + const result = await client.query({ + query: + "SELECT count() AS supported FROM system.settings WHERE name = 'vector_search_use_quantized_codes'", + format: "JSONEachRow", + }); + const rows = await result.json<{ supported: string }>(); + const quantized = + process.env.RAG_V2_QUANTIZED !== "false" && + Number(rows[0]?.supported) > 0 && + provider.dimensions % 8 === 0; + await migrate(client, provider.dimensions, quantized); + console.log( + JSON.stringify({ + event: "ragv2.migrated", + dimensions: provider.dimensions, + quantized, + }), + ); +} catch { + console.error("RAG_V2_MIGRATION_FAILED"); + process.exitCode = 1; +} finally { + await client.close(); +} diff --git a/v2/src/pipeline.ts b/v2/src/pipeline.ts new file mode 100644 index 00000000..222dfe9b --- /dev/null +++ b/v2/src/pipeline.ts @@ -0,0 +1,275 @@ +import { createHash, randomBytes } from "node:crypto"; +import { chunks } from "./chunking"; +import { + RagError, + narrowScope, + type Document, + type IngestInput, + type Scope, + type Store, + type StoredChunk, +} from "./contracts"; +import { DocumentLocks, Gate } from "./limits"; +import { validateVectors, type EmbeddingProvider } from "./embeddings"; + +export class Pipeline { + private readonly locks = new DocumentLocks(); + private readonly admission = new Gate(2); + private lastVersion = 0n; + constructor( + private readonly store: Store, + private readonly provider: EmbeddingProvider, + ) {} + private version(floor = "0"): string { + if (BigInt(floor) > this.lastVersion) this.lastVersion = BigInt(floor); + const now = BigInt(Date.now()) * 1000000n; + this.lastVersion = now > this.lastVersion ? now : this.lastVersion + 1n; + return this.lastVersion.toString(); + } + private operation( + scope: Scope, + fileId: string, + subject: string, + key: string, + operation: string, + ): string { + return createHash("sha256") + .update( + JSON.stringify([ + scope.tenantId, + scope.namespaceId, + fileId, + subject, + key, + operation, + ]), + ) + .digest("hex"); + } + private async replay( + scope: Scope, + fileId: string, + operationKey: string, + requestHash: string, + signal: AbortSignal, + ): Promise<{ replay: Document | null; current: Document | null }> { + const receipt = await this.store.receipt(scope, operationKey, signal); + if (receipt) { + if (receipt.requestHash !== requestHash) + throw new RagError("IDEMPOTENCY_CONFLICT", 409); + return { replay: receipt.document, current: null }; + } + // Recover a lost acknowledgment between publication and the receipt insert. + const current = await this.store.get(scope, fileId, signal); + if (current?.operationKey === operationKey) { + if (current.requestHash !== requestHash) + throw new RagError("IDEMPOTENCY_CONFLICT", 409); + await this.store.record( + scope, + operationKey, + { requestHash, document: current }, + signal, + ); + return { replay: current, current }; + } + if (current) { + // A later revision must not overwrite the only receipt of a published write. + await this.store.record( + scope, + current.operationKey, + { requestHash: current.requestHash, document: current }, + signal, + ); + } + return { replay: null, current }; + } + async ingest( + scope: Scope, + fileId: string, + subject: string, + key: string, + input: IngestInput, + signal: AbortSignal, + provenance: { + actorKind: "user" | "service"; + sourceClass: "asserted" | "extracted"; + } = { actorKind: "user", sourceClass: "asserted" }, + ): Promise { + scope = narrowScope(scope, fileId); + return this.admission.run(signal, () => + this.locks.run( + JSON.stringify([scope.tenantId, scope.namespaceId, fileId]), + async () => { + const operationKey = this.operation( + scope, + fileId, + `${provenance.actorKind}:${subject}`, + key, + "put", + ); + const requestHash = createHash("sha256") + .update( + JSON.stringify([ + input, + provenance.sourceClass, + this.provider.spaceId, + "chunk-1500-overlap150-v1", + ]), + ) + .digest("hex"); + const { replay, current } = await this.replay( + scope, + fileId, + operationKey, + requestHash, + signal, + ); + if (replay) return replay; + if (input.ifMatch && current?.generation !== input.ifMatch) + throw new RagError("PRECONDITION_FAILED", 409); + const generation = randomBytes(16).toString("hex"); + let chunkCount = 0; + let batchIndex = 0; + let batch: ReturnType extends Generator + ? T[] + : never = []; + let bytes = 0; + let pendingInsert = Promise.resolve(); + let insertionError: unknown; + const flush = async () => { + if (!batch.length) return; + signal.throwIfAborted(); + if (insertionError) throw insertionError; + const embeddings = await this.provider.embedDocuments( + batch.map((chunk) => + [input.title, chunk.section.join(" > "), chunk.text] + .filter(Boolean) + .join("\n\n"), + ), + signal, + ); + validateVectors(embeddings, batch.length, this.provider.dimensions); + const rows: StoredChunk[] = batch.map((chunk, index) => ({ + ...chunk, + tenantId: scope.tenantId, + namespaceId: scope.namespaceId, + fileId, + generation, + spaceId: this.provider.spaceId, + actor: `${provenance.actorKind}:${subject}`, + sourceClass: provenance.sourceClass, + embedding: embeddings[index]!, + })); + await pendingInsert; + if (insertionError) throw insertionError; + pendingInsert = this.store + .insert(rows, `${generation}:${batchIndex++}`, signal) + .catch((error) => { + insertionError = error; + }); + chunkCount += batch.length; + batch = []; + bytes = 0; + }; + try { + for (const chunk of chunks(input.segments)) { + const rowBytes = + Buffer.byteLength(chunk.text) + + this.provider.dimensions * 24 + + 2048; + if ( + batch.length && + (batch.length >= 32 || bytes + rowBytes > 1024 * 1024) + ) + await flush(); + batch.push(chunk); + bytes += rowBytes; + } + await flush(); + } finally { + await pendingInsert; + } + if (insertionError) throw insertionError; + if (!chunkCount) throw new RagError("EMPTY_DOCUMENT", 422); + signal.throwIfAborted(); + const document: Document = { + tenantId: scope.tenantId, + namespaceId: scope.namespaceId, + fileId, + generation, + version: this.version(current?.version), + state: "ready", + title: input.title, + original: input.original, + actor: `${provenance.actorKind}:${subject}`, + sourceClass: provenance.sourceClass, + spaceId: this.provider.spaceId, + chunkCount, + operationKey, + requestHash, + }; + await this.store.publish(document, signal); + await this.store.record( + scope, + operationKey, + { requestHash, document }, + signal, + ); + return document; + }, + ), + ); + } + async delete( + scope: Scope, + fileId: string, + subject: string, + key: string, + signal: AbortSignal, + actorKind: "user" | "service" = "user", + ): Promise { + scope = narrowScope(scope, fileId); + return this.admission.run(signal, () => + this.locks.run( + JSON.stringify([scope.tenantId, scope.namespaceId, fileId]), + async () => { + const operationKey = this.operation( + scope, + fileId, + `${actorKind}:${subject}`, + key, + "delete", + ); + const requestHash = createHash("sha256") + .update("delete-v1") + .digest("hex"); + const { replay, current } = await this.replay( + scope, + fileId, + operationKey, + requestHash, + signal, + ); + if (replay) return replay; + if (!current) throw new RagError("NOT_FOUND", 404); + const document: Document = { + ...current, + state: "deleted", + actor: `${actorKind}:${subject}`, + version: this.version(current?.version), + operationKey, + requestHash, + }; + await this.store.publish(document, signal); + await this.store.record( + scope, + operationKey, + { requestHash, document }, + signal, + ); + return document; + }, + ), + ); + } +} diff --git a/v2/src/schema.ts b/v2/src/schema.ts new file mode 100644 index 00000000..1a3f3ca9 --- /dev/null +++ b/v2/src/schema.ts @@ -0,0 +1,45 @@ +import type { ClickHouseClient } from "@clickhouse/client"; + +export async function migrate( + client: ClickHouseClient, + dimensions: number, + quantized = false, +): Promise { + if ( + !Number.isInteger(dimensions) || + dimensions < 8 || + dimensions > 8192 || + (quantized && dimensions % 8) + ) + throw new Error("INVALID_DIMENSIONS"); + const statements = [ + `CREATE TABLE IF NOT EXISTS documents ( + tenant_id LowCardinality(String), namespace_id LowCardinality(String), file_id String, + generation FixedString(32), version UInt64, state Enum8('ready'=1,'deleted'=2), title String, + original String, actor String, source_class Enum8('asserted'=1,'extracted'=2), space_id LowCardinality(String), chunk_count UInt32, operation_key FixedString(64), request_hash FixedString(64) + ) ENGINE = ReplacingMergeTree(version) ORDER BY (tenant_id, namespace_id, file_id) SETTINGS index_granularity=64, non_replicated_deduplication_window=10000`, + `CREATE TABLE IF NOT EXISTS chunks ( + tenant_id LowCardinality(String), namespace_id LowCardinality(String), space_id LowCardinality(String), + file_id String, generation FixedString(32), chunk_index UInt32, content String, + content_lc String MATERIALIZED lowerUTF8(content), page Nullable(UInt16), segment UInt16, + char_start UInt32, char_end UInt32, section Array(String), actor String, source_class Enum8('asserted'=1,'extracted'=2), + embedding Array(${quantized ? "BFloat16" : "Float32"}) ${quantized ? `CODEC(Quantized('rabitq', ${dimensions}, 0))` : "CODEC(NONE)"}, + CONSTRAINT vector_width CHECK length(embedding) = ${dimensions}, + INDEX content_text content_lc TYPE text(tokenizer='splitByNonAlpha') GRANULARITY 100000000 + ) ENGINE=MergeTree PARTITION BY cityHash64(tenant_id, namespace_id)%16 + ORDER BY (tenant_id, namespace_id, space_id, file_id, generation, chunk_index) + SETTINGS non_replicated_deduplication_window=10000, exclude_materialize_skip_indexes_on_merge='content_text'`, + `CREATE TABLE IF NOT EXISTS operations ( + tenant_id LowCardinality(String), namespace_id LowCardinality(String), operation_key FixedString(64), receipt String, version UInt64 + ) ENGINE=ReplacingMergeTree(version) ORDER BY (tenant_id, namespace_id, operation_key) + SETTINGS index_granularity=64, non_replicated_deduplication_window=10000`, + ]; + for (const query of statements) + await client.command({ + query, + clickhouse_settings: { + log_comment: "ragv2.migrate", + ...(quantized ? { enable_quantized_codec: "1" } : {}), + }, + }); +} diff --git a/v2/test/app.test.ts b/v2/test/app.test.ts new file mode 100644 index 00000000..8d90fe16 --- /dev/null +++ b/v2/test/app.test.ts @@ -0,0 +1,206 @@ +import { expect, test } from "bun:test"; +import { createApp } from "../src/app"; +import { auth, MemoryStore, provider, token } from "./helpers"; + +const path = "/v2/namespaces/library-a/documents/file-a"; +const body = { + segments: [{ kind: "page", index: 1, text: "original alpha document" }], + original: { + fileId: "file-a", + revision: "revision-a", + sha256: "a".repeat(64), + filename: "original.pdf", + mediaType: "application/pdf", + }, +}; +async function headers(overrides = {}) { + return { + Authorization: `Bearer ${await token(overrides)}`, + "Content-Type": "application/json", + "Idempotency-Key": "ingest-a", + }; +} + +test("authenticated ingest, search, ordered context, original revision, and deletion", async () => { + const store = new MemoryStore(); + const app = createApp({ auth, store, provider, writer: true }); + const signed = await headers(); + expect( + ( + await app.request(path, { + method: "PUT", + headers: signed, + body: JSON.stringify(body), + }) + ).status, + ).toBe(200); + const response = await app.request("/v2/search", { + method: "POST", + headers: signed, + body: JSON.stringify({ + query: "alpha", + namespaces: [{ namespaceId: "library-a" }], + }), + }); + const result = await response.json(); + expect(result.hits[0].original).toEqual(body.original); + expect(result.hits[0].page).toBe(1); + expect(result.hits[0].start).toBe(0); + expect( + (await app.request(`${path}/context`, { headers: signed })).status, + ).toBe(200); + expect( + ( + await app.request(path, { + method: "DELETE", + headers: { ...signed, "Idempotency-Key": "delete-a" }, + }) + ).status, + ).toBe(204); + expect((await app.request(path, { headers: signed })).status).toBe(404); +}); +test("unauthenticated and unauthorized writes fail before body consumption or I/O", async () => { + const store = new MemoryStore(); + const app = createApp({ auth, store, provider, writer: true }); + expect( + (await app.request(path, { method: "PUT", body: "not json" })).status, + ).toBe(401); + const signed = await headers({ + grants: [ + { + namespaceId: "library-a", + resourceKind: "document", + operations: ["read"], + resourceIds: ["file-b"], + }, + ], + }); + expect( + ( + await app.request(path, { + method: "PUT", + headers: signed, + body: "not json", + }) + ).status, + ).toBe(403); + expect(store.calls).toBe(0); +}); +test("body fields cannot supply tenant or widen entity scope; API replicas refuse writes", async () => { + const store = new MemoryStore(); + const signed = await headers(); + const app = createApp({ auth, store, provider, writer: true }); + for (const extra of [{ tenantId: "foreign" }, { entity_id: "victim" }]) + expect( + ( + await app.request(path, { + method: "PUT", + headers: signed, + body: JSON.stringify({ ...body, ...extra }), + }) + ).status, + ).toBe(400); + expect( + ( + await createApp({ auth, store, provider }).request(path, { + method: "PUT", + headers: signed, + body: JSON.stringify(body), + }) + ).status, + ).toBe(503); + expect(store.calls).toBe(0); +}); +test("search grants restrict files and tenants even when the caller omits file filters", async () => { + const store = new MemoryStore(); + const app = createApp({ auth, store, provider, writer: true }); + await app.request(path, { + method: "PUT", + headers: await headers(), + body: JSON.stringify(body), + }); + const signed = await headers({ + grants: [ + { + namespaceId: "library-a", + resourceKind: "document", + operations: ["read"], + resourceIds: ["file-b"], + }, + ], + }); + const request = { + query: "alpha", + namespaces: [{ namespaceId: "library-a" }], + }; + expect( + ( + await ( + await app.request("/v2/search", { + method: "POST", + headers: signed, + body: JSON.stringify(request), + }) + ).json() + ).hits, + ).toEqual([]); + const foreign = await headers({ tenant_id: "tenant-b" }); + expect( + ( + await ( + await app.request("/v2/search", { + method: "POST", + headers: foreign, + body: JSON.stringify(request), + }) + ).json() + ).hits, + ).toEqual([]); +}); +test("raw originals use a separate extraction adapter, bind the digest, and publish only validated text", async () => { + const store = new MemoryStore(); + let calls = 0; + const app = createApp({ + auth, + store, + provider, + writer: true, + extractor: { + async extract() { + calls++; + return [{ kind: "page", index: 1, text: "extracted alpha" }]; + }, + }, + }); + const bytes = new TextEncoder().encode("%PDF fixture"); + const digest = new Bun.CryptoHasher("sha256").update(bytes).digest("hex"); + const original = { ...body.original, sha256: digest }; + const signed = { + ...(await headers()), + "Content-Type": "application/octet-stream", + "Content-Length": String(bytes.length), + "X-Rag-Format": "pdf", + "X-Rag-Original": JSON.stringify(original), + }; + expect( + ( + await app.request(`${path}/content`, { + method: "PUT", + headers: signed, + body: bytes, + }) + ).status, + ).toBe(200); + expect(calls).toBe(1); + expect(store.chunks[0]!.text).toBe("extracted alpha"); + expect( + ( + await app.request(`${path}/content`, { + method: "PUT", + headers: { ...signed, "X-Rag-Original": JSON.stringify(body.original) }, + body: bytes, + }) + ).status, + ).toBe(400); + expect(calls).toBe(1); +}); diff --git a/v2/test/auth.test.ts b/v2/test/auth.test.ts new file mode 100644 index 00000000..917459b9 --- /dev/null +++ b/v2/test/auth.test.ts @@ -0,0 +1,68 @@ +import { expect, test } from "bun:test"; +import { authorize, createAuthenticator } from "../src/auth"; +import { auth, token } from "./helpers"; + +test("verifies short-lived asymmetric service identities, not browser/code-api tokens", async () => { + expect((await auth.verify(`Bearer ${await token()}`)).tenant_id).toBe( + "tenant-a", + ); + for (const override of [ + { aud: "codeapi" }, + { iss: "foreign" }, + { tenant_id: undefined }, + { grants: undefined }, + { exp: 1 }, + { exp: Math.floor(Date.now() / 1000) + 301 }, + { nbf: Math.floor(Date.now() / 1000) + 600 }, + { iat: Math.floor(Date.now() / 1000) + 600 }, + { sub: "" }, + ]) { + await expect( + auth.verify(`Bearer ${await token(override)}`), + ).rejects.toMatchObject({ code: "AUTH_INVALID" }); + } + await expect(auth.verify(undefined)).rejects.toMatchObject({ status: 401 }); + const signed = await token(); + await expect( + auth.verify(`Bearer ${signed.slice(0, -8)}xxxxxxxx`), + ).rejects.toMatchObject({ status: 401 }); +}); +test("scope cannot be widened by file ids, operation, library, or tenant assertions", async () => { + const principal = await auth.verify( + `Bearer ${await token({ grants: [{ namespaceId: "library-a", resourceKind: "document", operations: ["read"], resourceIds: ["file-a"] }] })}`, + ); + expect(authorize(principal, "library-a", "read").resourceIds).toEqual([ + "file-a", + ]); + expect(() => authorize(principal, "library-a", "read", ["file-b"])).toThrow( + "SCOPE_DENIED", + ); + expect(() => authorize(principal, "library-a", "delete", ["file-a"])).toThrow( + "SCOPE_DENIED", + ); + expect(() => authorize(principal, "library-b", "read")).toThrow( + "SCOPE_DENIED", + ); +}); +test("fails closed on missing, private, symmetric, or ambiguous verification keys", async () => { + for (const keys of [ + [], + [{ kid: "x", alg: "HS256", kty: "oct", k: "secret" }], + [{ kid: "x", alg: "EdDSA", kty: "OKP", crv: "Ed25519", d: "private" }], + ]) { + await expect(createAuthenticator({ keys })).rejects.toThrow( + "INVALID_VERIFICATION_KEYS", + ); + } +}); + +test("resource kinds and actor kinds cannot be forged or promoted by request fields", async () => { + await expect( + auth.verify(`Bearer ${await token({ actor_kind: "admin" })}`), + ).rejects.toMatchObject({ code: "AUTH_INVALID" }); + await expect( + auth.verify( + `Bearer ${await token({ grants: [{ namespaceId: "library-a", resourceKind: "memory", operations: ["read"] }] })}`, + ), + ).rejects.toMatchObject({ code: "AUTH_INVALID" }); +}); diff --git a/v2/test/chunking.test.ts b/v2/test/chunking.test.ts new file mode 100644 index 00000000..2fdf6eae --- /dev/null +++ b/v2/test/chunking.test.ts @@ -0,0 +1,77 @@ +import { expect, test } from "bun:test"; +import { chunks } from "../src/chunking"; +import { ingestSchema } from "../src/contracts"; +import { Gate, readBounded } from "../src/limits"; +import { signal } from "./helpers"; + +test("chunks never cross page boundaries and preserve UTF-16 source offsets and section breadcrumbs", () => { + const segments = [ + { + kind: "page" as const, + index: 1, + text: "# Heading\n" + "πŸ˜€word ".repeat(500), + }, + { kind: "page" as const, index: 2, text: "Second page" }, + ]; + const result = [...chunks(segments)]; + expect(result.map((chunk) => chunk.index)).toEqual( + result.map((_, index) => index), + ); + for (const chunk of result) { + expect(chunk.text).toBe( + segments[chunk.page! - 1]!.text.slice(chunk.start, chunk.end), + ); + expect(chunk.text.isWellFormed()).toBe(true); + } + expect(result[0]!.section).toEqual(["Heading"]); + expect(result.at(-1)!.page).toBe(2); +}); +test("rejects unordered, oversized and empty extraction output", () => { + for (const segments of [ + [{ kind: "page", index: 2, text: "text" }], + [{ kind: "document", index: 1, text: " " }], + [{ kind: "document", index: 1, text: "x".repeat(1024 * 1024 + 1) }], + ]) + expect(() => ingestSchema.parse({ segments })).toThrow(); +}); +test("admission is bounded, queued cancellation frees its position, and slots recover after failure", async () => { + const gate = new Gate(1, 1); + let release: () => void = () => {}; + const active = gate.run( + signal(), + () => + new Promise((resolve) => { + release = resolve; + }), + ); + const controller = new AbortController(); + const queued = gate.run(controller.signal, async () => {}); + await expect(gate.run(signal(), async () => {})).rejects.toMatchObject({ + code: "OVERLOADED", + }); + controller.abort(); + await expect(queued).rejects.toThrow(); + release(); + await active; + await expect( + gate.run(signal(), async () => { + throw Error("work"); + }), + ).rejects.toThrow(); + expect(await gate.run(signal(), async () => 1)).toBe(1); +}); +test("body cap rejects without reading the rest of an untrusted stream", async () => { + let cancelled = false; + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new Uint8Array(20)); + }, + cancel() { + cancelled = true; + }, + }); + await expect(readBounded(body, 10, signal())).rejects.toMatchObject({ + code: "BODY_LIMIT", + }); + expect(cancelled).toBe(true); +}); diff --git a/v2/test/clickhouse.test.ts b/v2/test/clickhouse.test.ts new file mode 100644 index 00000000..f2ceca61 --- /dev/null +++ b/v2/test/clickhouse.test.ts @@ -0,0 +1,44 @@ +import { expect, test } from "bun:test"; +import { scopeWhere, liveWhere, vectorSQL } from "../src/clickhouse"; + +test("library/file predicates preserve correlated grants instead of cross-product access", () => { + const where = scopeWhere([ + { tenantId: "tenant-a", namespaceId: "a", resourceIds: ["one"] }, + { tenantId: "tenant-a", namespaceId: "b", resourceIds: ["two"] }, + ]); + expect(where.sql).toContain( + "namespace_id = {lib0:String} AND file_id IN {files0:Array(String)}", + ); + expect(where.sql).toContain( + "namespace_id = {lib1:String} AND file_id IN {files1:Array(String)}", + ); + expect(where.params).toEqual({ + tenant: "tenant-a", + lib0: "a", + files0: ["one"], + lib1: "b", + files1: ["two"], + }); + expect(() => + scopeWhere([ + { tenantId: "a", namespaceId: "x" }, + { tenantId: "b", namespaceId: "x" }, + ]), + ).toThrow(); +}); +test("publication is a scoped semi-join filter; the outer vector scan stays eligible", () => { + const where = liveWhere( + [{ tenantId: "a", namespaceId: "x", resourceIds: ["one"] }], + "space", + ); + const sql = vectorSQL(where.sql); + expect(sql).toContain("FROM documents FINAL"); + expect(sql).toContain("state = 'ready'"); + expect(sql).toContain( + "cosineDistance(embedding, {vector:Array(Float32)}) AS distance", + ); + expect(sql).toContain("ORDER BY distance ASC LIMIT {k:UInt32}"); + expect(sql).not.toContain("FROM chunks FINAL"); + expect(sql).not.toContain(" JOIN "); + expect(sql).not.toContain("LIMIT 1 BY"); +}); diff --git a/v2/test/embeddings.test.ts b/v2/test/embeddings.test.ts new file mode 100644 index 00000000..5e6fa9d0 --- /dev/null +++ b/v2/test/embeddings.test.ts @@ -0,0 +1,105 @@ +import { expect, test } from "bun:test"; +import { openAICompatible, QueryEmbeddings } from "../src/embeddings"; +import { provider, signal } from "./helpers"; + +async function serving( + handler: (request: Request) => Response | Promise, + work: (url: string) => Promise, +) { + const server = Bun.serve({ hostname: "127.0.0.1", port: 0, fetch: handler }); + try { + await work(`http://127.0.0.1:${server.port}/embeddings`); + } finally { + await server.stop(true); + } +} +test("provider reconstructs reordered vectors by index and binds dimensions to its space", async () => { + await serving( + () => + Response.json({ + data: [ + { index: 1, embedding: [0, 1, 0, 0, 0, 0, 0, 0] }, + { index: 0, embedding: [1, 0, 0, 0, 0, 0, 0, 0] }, + ], + }), + async (endpoint) => { + const adapter = openAICompatible({ + endpoint, + apiKey: "test", + model: "model", + dimensions: 8, + }); + const result = await adapter.embedDocuments(["one", "two"], signal()); + expect(result[0]![0]).toBe(1); + expect(result[1]![1]).toBe(1); + expect(adapter.spaceId).not.toBe( + openAICompatible({ + endpoint, + apiKey: "test", + model: "model", + dimensions: 16, + }).spaceId, + ); + }, + ); +}); +test("malformed vectors, duplicate indices, and redirects fail closed", async () => { + for (const data of [ + [{ index: 0, embedding: [1] }], + [ + { index: 0, embedding: Array(8).fill(1) }, + { index: 0, embedding: Array(8).fill(1) }, + ], + ]) { + await serving( + () => Response.json({ data }), + async (endpoint) => { + await expect( + openAICompatible({ + endpoint, + apiKey: "test", + model: "model", + dimensions: 8, + }).embedDocuments(["one", "two"], signal()), + ).rejects.toMatchObject({ code: "INVALID_EMBEDDINGS" }); + }, + ); + } + let requests = 0; + await serving( + () => { + requests++; + return Response.redirect("http://127.0.0.1:1/secret"); + }, + async (endpoint) => { + await expect( + openAICompatible({ + endpoint, + apiKey: "test", + model: "model", + dimensions: 8, + }).embedQuery("one", signal()), + ).rejects.toThrow(); + }, + ); + expect(requests).toBeLessThanOrEqual(2); +}); +test("query single-flight isolates waiter cancellation and failures do not poison the cache", async () => { + let count = 0; + const cached = new QueryEmbeddings({ + ...provider, + async embedQuery(text, sig) { + count++; + await Bun.sleep(20); + return provider.embedQuery(text, sig); + }, + }); + const controller = new AbortController(); + const a = cached.get("same query", controller.signal); + const b = cached.get("same query", signal()); + controller.abort(); + await expect(a).rejects.toThrow(); + await b; + await cached.get("same query", signal()); + expect(count).toBe(1); +}); diff --git a/v2/test/extraction.test.ts b/v2/test/extraction.test.ts new file mode 100644 index 00000000..01d35ba7 --- /dev/null +++ b/v2/test/extraction.test.ts @@ -0,0 +1,59 @@ +import { expect, test } from "bun:test"; +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { createServer } from "node:http"; +import { once } from "node:events"; +import { unixExtractor } from "../src/extraction"; +import { signal } from "./helpers"; + +test("Bun extraction client uses exact binary length, validates digest and one-based pages, and rejects mismatches", async () => { + const dir = await mkdtemp(join(tmpdir(), "ragv2-extract-")); + const socket = join(dir, "service.sock"); + let wrong = false; + const server = createServer(async (req, res) => { + const parts: Buffer[] = []; + for await (const chunk of req) parts.push(chunk); + expect(req.headers["content-length"]).toBe("3"); + expect(req.headers["x-extraction-version"]).toBe("1"); + expect(Buffer.concat(parts).toString()).toBe("pdf"); + const result = { + version: 1, + operation: "document.extract-text", + format: "pdf", + inputSha256: wrong ? "b".repeat(64) : "a".repeat(64), + textBytes: 5, + segments: [{ kind: "page", index: 1, text: "alpha" }], + }; + res.writeHead(200, { "Content-Type": "application/json" }); + res.end(JSON.stringify(result)); + }); + server.listen(socket); + await once(server, "listening"); + try { + const extractor = unixExtractor(socket); + expect( + ( + await extractor.extract( + new TextEncoder().encode("pdf"), + "pdf", + "a".repeat(64), + signal(), + ) + )[0]!.index, + ).toBe(1); + wrong = true; + await expect( + extractor.extract( + new TextEncoder().encode("pdf"), + "pdf", + "a".repeat(64), + signal(), + ), + ).rejects.toMatchObject({ code: "INVALID_EXTRACTION_RESULT" }); + } finally { + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + await rm(dir, { recursive: true }); + } +}); diff --git a/v2/test/helpers.ts b/v2/test/helpers.ts new file mode 100644 index 00000000..e68a1e28 --- /dev/null +++ b/v2/test/helpers.ts @@ -0,0 +1,144 @@ +import { generateKeyPair, exportJWK, SignJWT } from "jose"; +import { createAuthenticator } from "../src/auth"; +import { deterministicProvider } from "../src/embeddings"; +import type { + Chunk, + Document, + Hit, + Receipt, + Scope, + SearchInput, + Store, + StoredChunk, +} from "../src/contracts"; + +export const signal = () => AbortSignal.timeout(5000); +export const provider = deterministicProvider(64); +export const keypair = await generateKeyPair("EdDSA"); +const publicKey = await exportJWK(keypair.publicKey); +export const auth = await createAuthenticator({ + keys: [{ ...publicKey, kid: "test-key", alg: "EdDSA" }], +}); +export async function token( + overrides: Record = {}, +): Promise { + const now = Math.floor(Date.now() / 1000); + return new SignJWT({ + tenant_id: "tenant-a", + grants: [ + { + namespaceId: "library-a", + resourceKind: "document", + operations: ["read", "write", "delete"], + }, + ], + iss: "librechat", + aud: "ragapi", + sub: "user-a", + actor_kind: "user", + iat: now, + nbf: now, + exp: now + 60, + jti: crypto.randomUUID(), + ...overrides, + }) + .setProtectedHeader({ alg: "EdDSA", typ: "JWT", kid: "test-key" }) + .sign(keypair.privateKey); +} +export class MemoryStore implements Store { + readonly documents = new Map(); + readonly chunks: StoredChunk[] = []; + readonly receipts = new Map(); + readonly events: string[] = []; + failBatch = 0; + failRecord = false; + calls = 0; + private docKey(scope: Scope, fileId: string) { + return JSON.stringify([scope.tenantId, scope.namespaceId, fileId]); + } + async get(scope: Scope, fileId: string): Promise { + this.calls++; + return this.documents.get(this.docKey(scope, fileId)) ?? null; + } + async receipt(scope: Scope, key: string): Promise { + this.calls++; + return this.receipts.get(this.docKey(scope, key)) ?? null; + } + async insert(chunks: readonly StoredChunk[]): Promise { + this.calls++; + this.events.push("insert"); + if ( + this.failBatch && + this.events.filter((event) => event === "insert").length === + this.failBatch + ) + throw new Error("insertion failed"); + this.chunks.push(...chunks); + } + async publish(document: Document): Promise { + this.calls++; + this.events.push("publish"); + this.documents.set(this.docKey(document, document.fileId), document); + } + async record(scope: Scope, key: string, receipt: Receipt): Promise { + this.calls++; + if (this.failRecord) throw new Error("lost acknowledgment"); + this.receipts.set(this.docKey(scope, key), receipt); + } + async context(scope: Scope, document: Document): Promise { + this.calls++; + return this.chunks.filter( + (chunk) => + chunk.tenantId === scope.tenantId && + chunk.namespaceId === scope.namespaceId && + chunk.fileId === document.fileId && + chunk.generation === document.generation, + ); + } + async search( + scopes: readonly Scope[], + vector: readonly number[], + spaceId: string, + input: SearchInput, + ): Promise { + this.calls++; + return this.chunks + .flatMap((chunk) => { + const scope = scopes.find( + (scope) => + scope.tenantId === chunk.tenantId && + scope.namespaceId === chunk.namespaceId && + (scope.resourceIds === undefined || + scope.resourceIds.includes(chunk.fileId)), + ); + const doc = scope + ? this.documents.get(this.docKey(scope, chunk.fileId)) + : null; + if ( + !doc || + doc.state !== "ready" || + doc.generation !== chunk.generation || + chunk.spaceId !== spaceId + ) + return []; + const distance = cosine(chunk.embedding, vector); + return [ + { ...chunk, distance, score: 1 - distance, original: doc.original }, + ]; + }) + .sort((a, b) => a.distance - b.distance) + .slice(0, input.k); + } + async close(): Promise {} +} +function cosine(a: readonly number[], b: readonly number[]): number { + let dot = 0; + let na = 0; + let nb = 0; + for (let i = 0; i < a.length; i++) { + dot += a[i]! * b[i]!; + na += a[i]! ** 2; + nb += b[i]! ** 2; + } + return 1 - dot / Math.sqrt(na * nb); +} diff --git a/v2/test/integration.test.ts b/v2/test/integration.test.ts new file mode 100644 index 00000000..c3b65535 --- /dev/null +++ b/v2/test/integration.test.ts @@ -0,0 +1,242 @@ +import { expect, test } from "bun:test"; +import { + ClickHouseStore, + databaseClient, + liveWhere, + vectorSQL, +} from "../src/clickhouse"; +import { migrate } from "../src/schema"; +import { Pipeline } from "../src/pipeline"; +import { ingestSchema, searchSchema } from "../src/contracts"; +import { createApp } from "../src/app"; +import { auth, provider, signal, token } from "./helpers"; + +const url = process.env.RAG_V2_TEST_CLICKHOUSE_URL; +for (const quantized of [false, true]) { + test.skipIf(!url)( + `real ClickHouse ${quantized ? "quantized" : "plain"}: publication, retries, scope, hybrid, originals, and deletion`, + async () => { + const database = `ragv2_test_${crypto.randomUUID().replaceAll("-", "")}`; + const config = { + url: url!, + database: "default", + username: "default", + password: "", + }; + const admin = databaseClient(config); + let store: ClickHouseStore | undefined; + try { + await admin.command({ query: `CREATE DATABASE ${database}` }); + const client = databaseClient({ ...config, database }); + store = new ClickHouseStore(client, database, 64); + await migrate(client, 64, quantized); + await store.probe(); + expect(store.quantizedCodesEnabled).toBe(quantized); + const scope = { tenantId: "tenant-a", namespaceId: "library-a" }; + const pipeline = new Pipeline(store, provider); + const input = ingestSchema.parse({ + title: "Alpha", + segments: [ + { kind: "page", index: 1, text: "alpha recovery document" }, + ], + original: { + fileId: "file-a", + revision: "one", + sha256: "a".repeat(64), + filename: "original.pdf", + mediaType: "application/pdf", + }, + }); + const doc = await pipeline.ingest( + scope, + "file-a", + "user-a", + "one", + input, + signal(), + ); + const vector = await provider.embedQuery("alpha recovery", signal()); + const search = searchSchema.parse({ + query: "alpha recovery", + namespaces: [{ namespaceId: "library-a" }], + precision: quantized ? "quantized" : "exact", + }); + let hits = await store.search( + [scope], + vector, + provider.spaceId, + search, + signal(), + ); + expect(hits).toHaveLength(1); + expect(hits[0]!.original).toEqual(input.original); + expect(hits[0]!.page).toBe(1); + expect((await store.get(scope, "file-a", signal()))!.version).toBe( + doc.version, + ); + expect((await store.context(scope, doc, signal()))[0]!.text).toBe( + input.segments[0]!.text, + ); + expect( + ( + await pipeline.ingest( + scope, + "file-a", + "user-a", + "one", + input, + signal(), + ) + ).generation, + ).toBe(doc.generation); + expect( + await store.search( + [{ ...scope, tenantId: "tenant-b" }], + vector, + provider.spaceId, + search, + signal(), + ), + ).toEqual([]); + expect( + await store.search( + [{ ...scope, resourceIds: ["foreign"] }], + vector, + provider.spaceId, + search, + signal(), + ), + ).toEqual([]); + expect( + await store.search( + [scope], + vector, + "different-model", + search, + signal(), + ), + ).toEqual([]); + // A successfully inserted but unpublished batch must never be searchable. + const row = { + tenantId: scope.tenantId, + namespaceId: scope.namespaceId, + fileId: "staged", + generation: "b".repeat(32), + spaceId: provider.spaceId, + index: 0, + text: "alpha recovery staged", + page: 1, + segment: 1, + start: 0, + end: 21, + section: [], + actor: "user:user-a", + sourceClass: "asserted" as const, + embedding: vector, + }; + await store.insert([row], "retry-same-batch", signal()); + await store.insert([row], "retry-same-batch", signal()); + const count = await store.rows<{ n: string }>( + "SELECT count() AS n FROM chunks WHERE file_id = 'staged'", + {}, + signal(), + ); + expect(Number(count[0]!.n)).toBe(1); + expect( + ( + await store.search( + [scope], + vector, + provider.spaceId, + search, + signal(), + ) + ).every((hit) => hit.fileId !== "staged"), + ).toBe(true); + hits = await store.search( + [scope], + vector, + provider.spaceId, + { ...search, mode: "hybrid" }, + signal(), + ); + expect(hits).toHaveLength(1); + const where = liveWhere([scope], provider.spaceId); + const plan = await client.query({ + query: `EXPLAIN PLAN ${vectorSQL(where.sql)}`, + query_params: { ...where.params, vector, k: 5 }, + format: "TabSeparated", + clickhouse_settings: { + ...(quantized + ? { + vector_search_use_quantized_codes: "1", + vector_search_index_fetch_multiplier: "3", + } + : {}), + }, + }); + const text = await plan.text(); + console.log( + JSON.stringify({ + qualification: "vector-plan", + quantized, + plan: text, + }), + ); + expect(text).toContain("ReadFromMergeTree"); + if (quantized) expect(text).toMatch(/Quantized|Rescor|Lazy/i); + const updated = await pipeline.ingest( + scope, + "file-a", + "user-a", + "two", + { + ...input, + segments: [ + { kind: "page", index: 1, text: "replacement beta document" }, + ], + }, + signal(), + ); + expect( + ( + await store.search( + [scope], + vector, + provider.spaceId, + search, + signal(), + ) + ).every((hit) => hit.generation === updated.generation), + ).toBe(true); + const app = createApp({ auth, store, provider }); + const response = await app.request("/v2/search", { + method: "POST", + headers: { + Authorization: `Bearer ${await token()}`, + "Content-Type": "application/json", + }, + body: JSON.stringify(search), + }); + expect(response.status).toBe(200); + await pipeline.delete(scope, "file-a", "user-a", "delete", signal()); + expect( + await store.search( + [scope], + vector, + provider.spaceId, + search, + signal(), + ), + ).toEqual([]); + } finally { + await store?.close(); + await admin.command({ + query: `DROP DATABASE IF EXISTS ${database} SYNC`, + }); + await admin.close(); + } + }, + 60000, + ); +} diff --git a/v2/test/pipeline.test.ts b/v2/test/pipeline.test.ts new file mode 100644 index 00000000..1f04d2a1 --- /dev/null +++ b/v2/test/pipeline.test.ts @@ -0,0 +1,261 @@ +import { expect, test } from "bun:test"; +import { Pipeline } from "../src/pipeline"; +import { ingestSchema, searchSchema } from "../src/contracts"; +import { MemoryStore, provider, signal } from "./helpers"; + +const scope = { tenantId: "tenant-a", namespaceId: "library-a" }; +const input = (text = "alpha document") => + ingestSchema.parse({ + title: "Example", + segments: [{ kind: "page", index: 1, text }], + }); + +test("publishes only after complete bounded batches and carries ordered source locations", async () => { + const store = new MemoryStore(); + const pipeline = new Pipeline(store, provider); + const doc = await pipeline.ingest( + scope, + "file-a", + "user", + "key", + input("alpha ".repeat(18000)), + signal(), + ); + expect(store.events.at(-1)).toBe("publish"); + expect( + store.events.filter((event) => event === "insert").length, + ).toBeGreaterThan(1); + expect(doc.chunkCount).toBe(store.chunks.length); + expect(store.chunks.map((chunk) => chunk.index)).toEqual( + Array.from({ length: doc.chunkCount }, (_, index) => index), + ); + expect( + store.chunks.every( + (chunk) => chunk.page === 1 && chunk.embedding.length === 64, + ), + ).toBe(true); + expect(BigInt(doc.version)).toBeGreaterThan(BigInt(Number.MAX_SAFE_INTEGER)); +}); +test("partial replacement failure leaves the previous generation live and staged chunks invisible", async () => { + const store = new MemoryStore(); + const pipeline = new Pipeline(store, provider); + const old = await pipeline.ingest( + scope, + "file-a", + "user", + "old", + input(), + signal(), + ); + store.failBatch = 3; + await expect( + pipeline.ingest( + scope, + "file-a", + "user", + "new", + input("replacement ".repeat(12000)), + signal(), + ), + ).rejects.toThrow(); + expect((await store.get(scope, "file-a"))!.generation).toBe(old.generation); + const hits = await store.search( + [scope], + await provider.embedQuery("alpha", signal()), + provider.spaceId, + searchSchema.parse({ + query: "alpha", + namespaces: [{ namespaceId: "library-a" }], + }), + ); + expect(hits.every((hit) => hit.generation === old.generation)).toBe(true); +}); +test("durable idempotency prevents duplicate ingestion and rejects a reused key with another body", async () => { + const store = new MemoryStore(); + const pipeline = new Pipeline(store, provider); + const doc = await pipeline.ingest( + scope, + "file-a", + "user", + "key", + input(), + signal(), + ); + const replay = await new Pipeline(store, provider).ingest( + scope, + "file-a", + "user", + "key", + input(), + signal(), + ); + expect(replay).toEqual(doc); + expect(store.chunks).toHaveLength(1); + await expect( + pipeline.ingest( + scope, + "file-a", + "user", + "key", + input("different"), + signal(), + ), + ).rejects.toMatchObject({ code: "IDEMPOTENCY_CONFLICT" }); +}); +test("recovers publication acknowledgment loss without re-embedding", async () => { + const store = new MemoryStore(); + const pipeline = new Pipeline(store, provider); + store.failRecord = true; + await expect( + pipeline.ingest(scope, "file-a", "user", "key", input(), signal()), + ).rejects.toThrow(); + store.failRecord = false; + await new Pipeline(store, provider).ingest( + scope, + "file-a", + "user", + "key", + input(), + signal(), + ); + expect(store.chunks).toHaveLength(1); +}); +test("same-document writers and deletes cannot race, and failed preconditions do not insert", async () => { + const store = new MemoryStore(); + const pipeline = new Pipeline(store, provider); + const a = pipeline.ingest(scope, "file-a", "user", "one", input(), signal()); + await expect( + pipeline.delete(scope, "file-a", "user", "delete", signal()), + ).rejects.toMatchObject({ code: "DOCUMENT_BUSY" }); + const doc = await a; + await expect( + pipeline.ingest( + scope, + "file-a", + "user", + "two", + { ...input(), ifMatch: "f".repeat(32) }, + signal(), + ), + ).rejects.toMatchObject({ code: "PRECONDITION_FAILED" }); + expect(store.chunks).toHaveLength(1); + await pipeline.delete(scope, "file-a", "user", "delete", signal()); + expect((await store.get(scope, "file-a"))!.state).toBe("deleted"); + expect((await store.get(scope, "file-a"))!.generation).toBe(doc.generation); +}); +test("invalid or cancelled embedding work never publishes", async () => { + const store = new MemoryStore(); + const invalid = new Pipeline(store, { + ...provider, + async embedDocuments() { + return [[NaN]]; + }, + }); + await expect( + invalid.ingest(scope, "file-a", "user", "key", input(), signal()), + ).rejects.toMatchObject({ code: "INVALID_EMBEDDINGS" }); + expect(store.documents.size).toBe(0); + const controller = new AbortController(); + controller.abort(); + await expect( + new Pipeline(store, provider).ingest( + scope, + "file-a", + "user", + "key", + input(), + controller.signal, + ), + ).rejects.toThrow(); + expect(store.documents.size).toBe(0); +}); + +test("embeds the next bounded batch while the preceding insert is in flight", async () => { + const store = new MemoryStore(); + let release = () => {}; + let embedded = 0; + const insert = store.insert.bind(store); + store.insert = async (rows) => { + if (rows[0]!.index === 0) + await new Promise((resolve) => { + release = resolve; + }); + await insert(rows); + }; + const pipeline = new Pipeline(store, { + ...provider, + async embedDocuments(texts, sig) { + embedded++; + return provider.embedDocuments(texts, sig); + }, + }); + const task = pipeline.ingest( + scope, + "file-a", + "user", + "overlap", + input("alpha ".repeat(18000)), + signal(), + ); + for (let tries = 0; embedded < 2 && tries < 100; tries++) await Bun.sleep(1); + expect(embedded).toBe(2); + expect(store.documents.size).toBe(0); + release(); + await task; +}); +test("persisted version floors survive clock skew and provenance is server-owned", async () => { + const store = new MemoryStore(); + const original = await new Pipeline(store, provider).ingest( + scope, + "file-a", + "user", + "first", + input(), + signal(), + ); + original.version = "9999999999999999999"; + const next = await new Pipeline(store, provider).ingest( + scope, + "file-a", + "service-worker", + "second", + input("replacement"), + signal(), + { actorKind: "service", sourceClass: "extracted" }, + ); + expect(BigInt(next.version)).toBeGreaterThan(BigInt(original.version)); + expect(next.actor).toBe("service:service-worker"); + expect(next.sourceClass).toBe("extracted"); + expect(store.chunks.at(-1)!.actor).toBe(next.actor); +}); + +test("a lost receipt cannot make a late retry resurrect an older revision", async () => { + const store = new MemoryStore(); + const pipeline = new Pipeline(store, provider); + store.failRecord = true; + await expect( + pipeline.ingest(scope, "file-a", "user", "old", input("old"), signal()), + ).rejects.toThrow(); + const original = (await store.get(scope, "file-a"))!; + store.failRecord = false; + const updated = await pipeline.ingest( + scope, + "file-a", + "user", + "new", + input("new"), + signal(), + ); + const replay = await new Pipeline(store, provider).ingest( + scope, + "file-a", + "user", + "old", + input("old"), + signal(), + ); + expect(replay.generation).toBe(original.generation); + expect((await store.get(scope, "file-a"))!.generation).toBe( + updated.generation, + ); +}); diff --git a/v2/tsconfig.json b/v2/tsconfig.json new file mode 100644 index 00000000..c70798a8 --- /dev/null +++ b/v2/tsconfig.json @@ -0,0 +1,16 @@ +{ + "compilerOptions": { + "target": "ES2024", + "module": "ESNext", + "moduleResolution": "Bundler", + "lib": ["ES2024", "DOM", "DOM.Iterable"], + "types": ["bun"], + "strict": true, + "noUncheckedIndexedAccess": true, + "noUnusedLocals": true, + "noUnusedParameters": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": ["src/**/*.ts", "test/**/*.ts", "scripts/**/*.ts"] +} From 583045932946cac7b9b220c9aeddebe2de8aba72 Mon Sep 17 00:00:00 2001 From: Lia Date: Sun, 4 Oct 2026 18:18:44 +0000 Subject: [PATCH 2/3] fix: Preserve Publication Preconditions and Bound Complete Context --- v2/scripts/benchmark.ts | 1 + v2/src/chunking.ts | 31 ++++++++- v2/src/clickhouse.ts | 17 ++++- v2/src/contracts.ts | 2 + v2/src/embeddings.ts | 7 +- v2/src/pipeline.ts | 17 +++-- v2/test/chunking.test.ts | 35 ++++++++++ v2/test/clickhouse.test.ts | 63 +++++++++++++++++ v2/test/integration.test.ts | 37 ++++++++++ v2/test/pipeline.test.ts | 131 ++++++++++++++++++++++++++++++++++++ 10 files changed, 329 insertions(+), 12 deletions(-) diff --git a/v2/scripts/benchmark.ts b/v2/scripts/benchmark.ts index 2a7557b2..e5c4fa61 100644 --- a/v2/scripts/benchmark.ts +++ b/v2/scripts/benchmark.ts @@ -36,6 +36,7 @@ const spaceId = `synthetic-${dimensions}`; const provider: EmbeddingProvider = { spaceId, dimensions, + maxInputBytes: 8191, async embedQuery() { return reference; }, diff --git a/v2/src/chunking.ts b/v2/src/chunking.ts index f4c030b1..582ed94d 100644 --- a/v2/src/chunking.ts +++ b/v2/src/chunking.ts @@ -1,4 +1,9 @@ -import type { Chunk, Segment } from "./contracts"; +import { + MAX_DOCUMENT_CHUNKS, + RagError, + type Chunk, + type Segment, +} from "./contracts"; function boundary(text: string, position: number): number { const code = text.charCodeAt(position - 1); @@ -46,7 +51,9 @@ export function* chunks( Math.min(position + size, section.end), ); const text = segment.text.slice(position, end); - if (text.trim()) + if (text.trim()) { + if (index >= MAX_DOCUMENT_CHUNKS) + throw new RagError("DOCUMENT_CHUNK_LIMIT", 413); yield { index: index++, text, @@ -56,6 +63,7 @@ export function* chunks( end, section: section.path, }; + } if (end >= section.end) break; position = Math.max( position + 1, @@ -65,3 +73,22 @@ export function* chunks( } } } + +export function documentEmbeddingInput( + chunk: Chunk, + title: string, + maxInputBytes: number, +): string { + const textBytes = Buffer.byteLength(chunk.text); + if (!Number.isSafeInteger(maxInputBytes) || textBytes > maxInputBytes) { + throw new RagError("EMBEDDING_INPUT_LIMIT", 422); + } + const header = [title, chunk.section.join(" > ")] + .filter(Boolean) + .join("\n\n"); + const bytes = Buffer.from(header); + let end = Math.min(bytes.length, Math.max(0, maxInputBytes - textBytes - 2)); + while (end > 0 && end < bytes.length && (bytes[end]! & 0xc0) === 0x80) end--; + const prefix = bytes.subarray(0, end).toString("utf8").trimEnd(); + return prefix ? `${prefix}\n\n${chunk.text}` : chunk.text; +} diff --git a/v2/src/clickhouse.ts b/v2/src/clickhouse.ts index 7837bd96..e98dab5f 100644 --- a/v2/src/clickhouse.ts +++ b/v2/src/clickhouse.ts @@ -11,6 +11,7 @@ import { Gate } from "./limits"; import { RagError, narrowScope, + MAX_DOCUMENT_CHUNKS, type Chunk, type Document, type Hit, @@ -546,16 +547,28 @@ export class ClickHouseStore implements Store { document: Document, signal: AbortSignal, ): Promise { + if (document.chunkCount > MAX_DOCUMENT_CHUNKS) + throw new RagError("DOCUMENT_CHUNK_LIMIT", 413); const where = liveWhere( [narrowScope(scope, document.fileId)], document.spaceId, ); const rows = await this.rows( - `SELECT chunk_index, content, page, segment, char_start, char_end, section FROM chunks WHERE ${where.sql} AND generation = {generation:String} ORDER BY chunk_index ASC LIMIT 10000`, - { ...where.params, generation: document.generation }, + `SELECT chunk_index, content, page, segment, char_start, char_end, section FROM chunks WHERE ${where.sql} AND generation = {generation:String} ORDER BY chunk_index ASC LIMIT {limit:UInt32}`, + { + ...where.params, + generation: document.generation, + limit: MAX_DOCUMENT_CHUNKS + 1, + }, signal, { select_sequential_consistency: "1" }, ); + if ( + rows.length !== document.chunkCount || + rows.some((row, index) => row.chunk_index !== index) + ) { + throw new RagError("CONTEXT_INCOMPLETE", 503); + } return rows.map((row) => ({ index: row.chunk_index, text: row.content, diff --git a/v2/src/contracts.ts b/v2/src/contracts.ts index ce1c7d23..7142ec58 100644 --- a/v2/src/contracts.ts +++ b/v2/src/contracts.ts @@ -1,5 +1,7 @@ import { z } from "zod"; +export const MAX_DOCUMENT_CHUNKS = 10000; + export const id = z.string().regex(/^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$/); export const originalSchema = z .object({ diff --git a/v2/src/embeddings.ts b/v2/src/embeddings.ts index f6f2c53b..3332d45f 100644 --- a/v2/src/embeddings.ts +++ b/v2/src/embeddings.ts @@ -6,6 +6,7 @@ import { RagError } from "./contracts"; export interface EmbeddingProvider { readonly spaceId: string; readonly dimensions: number; + readonly maxInputBytes: number; embedQuery(text: string, signal: AbortSignal): Promise; embedDocuments( texts: readonly string[], @@ -61,7 +62,7 @@ export function openAICompatible(options: { url.href, options.model, options.dimensions, - "raw-query/title-section-document-v1", + "raw-query/title-section-document-budgeted-v2", ]), ) .digest("hex"); @@ -145,6 +146,7 @@ export function openAICompatible(options: { return { spaceId, dimensions: options.dimensions, + maxInputBytes: 8191, embedQuery: async (text, signal) => (await embed([text], signal))[0]!, embedDocuments: embed, }; @@ -162,8 +164,9 @@ export function deterministicProvider(dimensions = 64): EmbeddingProvider { return vector; }; return { - spaceId: `test-only-token-hash-${dimensions}-v1`, + spaceId: `test-only-token-hash-${dimensions}-v2`, dimensions, + maxInputBytes: 8191, async embedQuery(text, signal) { signal.throwIfAborted(); return embed(text); diff --git a/v2/src/pipeline.ts b/v2/src/pipeline.ts index 222dfe9b..f8362925 100644 --- a/v2/src/pipeline.ts +++ b/v2/src/pipeline.ts @@ -1,5 +1,5 @@ import { createHash, randomBytes } from "node:crypto"; -import { chunks } from "./chunking"; +import { chunks, documentEmbeddingInput } from "./chunking"; import { RagError, narrowScope, @@ -113,7 +113,7 @@ export class Pipeline { input, provenance.sourceClass, this.provider.spaceId, - "chunk-1500-overlap150-v1", + "chunk-1500-overlap150-budgeted-prefix-v2", ]), ) .digest("hex"); @@ -125,7 +125,10 @@ export class Pipeline { signal, ); if (replay) return replay; - if (input.ifMatch && current?.generation !== input.ifMatch) + if ( + input.ifMatch && + (current?.state !== "ready" || current.generation !== input.ifMatch) + ) throw new RagError("PRECONDITION_FAILED", 409); const generation = randomBytes(16).toString("hex"); let chunkCount = 0; @@ -142,9 +145,11 @@ export class Pipeline { if (insertionError) throw insertionError; const embeddings = await this.provider.embedDocuments( batch.map((chunk) => - [input.title, chunk.section.join(" > "), chunk.text] - .filter(Boolean) - .join("\n\n"), + documentEmbeddingInput( + chunk, + input.title, + this.provider.maxInputBytes, + ), ), signal, ); diff --git a/v2/test/chunking.test.ts b/v2/test/chunking.test.ts index 2fdf6eae..2b7e2ae8 100644 --- a/v2/test/chunking.test.ts +++ b/v2/test/chunking.test.ts @@ -75,3 +75,38 @@ test("body cap rejects without reading the rest of an untrusted stream", async ( }); expect(cancelled).toBe(true); }); + +test("the shared chunk ceiling accepts exactly 10,000 sections and rejects the next", async () => { + const { MAX_DOCUMENT_CHUNKS } = await import("../src/contracts"); + const sections = (count: number) => [ + { kind: "document" as const, index: 1, text: "# x\n".repeat(count) }, + ]; + const accepted = [...chunks(sections(MAX_DOCUMENT_CHUNKS))]; + expect(accepted).toHaveLength(MAX_DOCUMENT_CHUNKS); + expect(accepted.at(-1)!.index).toBe(MAX_DOCUMENT_CHUNKS - 1); + expect(() => [...chunks(sections(MAX_DOCUMENT_CHUNKS + 1))]).toThrow( + "DOCUMENT_CHUNK_LIMIT", + ); +}); + +test("embedding prefixes end on UTF-8 boundaries without dropping chunk content", async () => { + const { documentEmbeddingInput } = await import("../src/chunking"); + const chunk = { + index: 0, + text: "ζ­£ζ–‡πŸ˜€", + page: null, + segment: 1, + start: 0, + end: 4, + section: ["η« θŠ‚πŸ˜€"], + }; + for (let limit = Buffer.byteLength(chunk.text); limit < 40; limit++) { + const embedded = documentEmbeddingInput(chunk, "ζ ‡ι’˜πŸ˜€", limit); + expect(Buffer.byteLength(embedded)).toBeLessThanOrEqual(limit); + expect(embedded.endsWith(chunk.text)).toBe(true); + expect(embedded.isWellFormed()).toBe(true); + } + expect(() => documentEmbeddingInput(chunk, "", 1)).toThrow( + "EMBEDDING_INPUT_LIMIT", + ); +}); diff --git a/v2/test/clickhouse.test.ts b/v2/test/clickhouse.test.ts index f2ceca61..26810f28 100644 --- a/v2/test/clickhouse.test.ts +++ b/v2/test/clickhouse.test.ts @@ -42,3 +42,66 @@ test("publication is a scoped semi-join filter; the outer vector scan stays elig expect(sql).not.toContain(" JOIN "); expect(sql).not.toContain("LIMIT 1 BY"); }); + +test("context never silently returns a clipped, missing, or duplicate chunk sequence", async () => { + const { ClickHouseStore, databaseClient } = await import("../src/clickhouse"); + const { MAX_DOCUMENT_CHUNKS, ingestSchema } = + await import("../src/contracts"); + const { Pipeline } = await import("../src/pipeline"); + const { MemoryStore, provider, signal } = await import("./helpers"); + const scope = { tenantId: "tenant-a", namespaceId: "library-a" }; + const document = await new Pipeline(new MemoryStore(), provider).ingest( + scope, + "file-a", + "user", + "key", + ingestSchema.parse({ + segments: [{ kind: "document", index: 1, text: "alpha" }], + }), + signal(), + ); + const client = databaseClient({ + url: "http://127.0.0.1:1", + database: "test", + username: "default", + password: "", + }); + const store = new ClickHouseStore(client, "test", 64); + const row = { + chunk_index: 0, + content: "alpha", + page: null, + segment: 1, + char_start: 0, + char_end: 5, + section: [], + }; + let rows = [row]; + let calls = 0; + store.rows = async (query: string, params: Record) => { + calls++; + expect(query).toContain("LIMIT {limit:UInt32}"); + expect(params.limit).toBe(MAX_DOCUMENT_CHUNKS + 1); + return rows as T[]; + }; + try { + expect(await store.context(scope, document, signal())).toHaveLength(1); + for (const invalid of [[], [row, row], [{ ...row, chunk_index: 1 }]]) { + rows = invalid; + await expect( + store.context(scope, document, signal()), + ).rejects.toMatchObject({ code: "CONTEXT_INCOMPLETE" }); + } + const before = calls; + await expect( + store.context( + scope, + { ...document, chunkCount: MAX_DOCUMENT_CHUNKS + 1 }, + signal(), + ), + ).rejects.toMatchObject({ code: "DOCUMENT_CHUNK_LIMIT" }); + expect(calls).toBe(before); + } finally { + await store.close(); + } +}); diff --git a/v2/test/integration.test.ts b/v2/test/integration.test.ts index c3b65535..afcbeed5 100644 --- a/v2/test/integration.test.ts +++ b/v2/test/integration.test.ts @@ -220,6 +220,16 @@ for (const quantized of [false, true]) { }); expect(response.status).toBe(200); await pipeline.delete(scope, "file-a", "user-a", "delete", signal()); + await expect( + pipeline.ingest( + scope, + "file-a", + "user-a", + "stale-edit", + { ...input, ifMatch: updated.generation }, + signal(), + ), + ).rejects.toMatchObject({ code: "PRECONDITION_FAILED" }); expect( await store.search( [scope], @@ -229,6 +239,33 @@ for (const quantized of [false, true]) { signal(), ), ).toEqual([]); + if (!quantized) { + const { MAX_DOCUMENT_CHUNKS } = await import("../src/contracts"); + const complete = await pipeline.ingest( + scope, + "ceiling", + "user-a", + "ceiling", + ingestSchema.parse({ + segments: [ + { + kind: "document", + index: 1, + text: "# x\n".repeat(MAX_DOCUMENT_CHUNKS), + }, + ], + }), + AbortSignal.timeout(30000), + ); + const context = await store.context( + scope, + complete, + AbortSignal.timeout(10000), + ); + expect(complete.chunkCount).toBe(MAX_DOCUMENT_CHUNKS); + expect(context).toHaveLength(MAX_DOCUMENT_CHUNKS); + expect(context.at(-1)!.index).toBe(MAX_DOCUMENT_CHUNKS - 1); + } } finally { await store?.close(); await admin.command({ diff --git a/v2/test/pipeline.test.ts b/v2/test/pipeline.test.ts index 1f04d2a1..6b2f09e5 100644 --- a/v2/test/pipeline.test.ts +++ b/v2/test/pipeline.test.ts @@ -259,3 +259,134 @@ test("a lost receipt cannot make a late retry resurrect an older revision", asyn updated.generation, ); }); + +test("a stale ifMatch cannot resurrect a deleted generation, but deliberate recreation still works", async () => { + const store = new MemoryStore(); + const pipeline = new Pipeline(store, provider); + const original = await pipeline.ingest( + scope, + "file-a", + "user", + "original", + input(), + signal(), + ); + await pipeline.delete(scope, "file-a", "user", "delete", signal()); + const inserts = store.events.filter((event) => event === "insert").length; + await expect( + new Pipeline(store, provider).ingest( + scope, + "file-a", + "user", + "stale-edit", + { ...input("changed"), ifMatch: original.generation }, + signal(), + ), + ).rejects.toMatchObject({ code: "PRECONDITION_FAILED" }); + expect(store.events.filter((event) => event === "insert")).toHaveLength( + inserts, + ); + expect((await store.get(scope, "file-a"))!.state).toBe("deleted"); + const recreated = await pipeline.ingest( + scope, + "file-a", + "user", + "recreate", + input("new"), + signal(), + ); + expect(recreated.state).toBe("ready"); + expect(recreated.generation).not.toBe(original.generation); + const updated = await pipeline.ingest( + scope, + "file-a", + "user", + "live-edit", + { ...input("edited"), ifMatch: recreated.generation }, + signal(), + ); + expect(updated.generation).not.toBe(recreated.generation); +}); + +test("excess Markdown sections cannot publish a context that would be truncated", async () => { + const store = new MemoryStore(); + const pipeline = new Pipeline(store, provider); + const original = await pipeline.ingest( + scope, + "file-a", + "user", + "original", + input(), + signal(), + ); + await expect( + pipeline.ingest( + scope, + "file-a", + "user", + "too-many", + input("# x\n".repeat(10001)), + signal(), + ), + ).rejects.toMatchObject({ code: "DOCUMENT_CHUNK_LIMIT", status: 413 }); + expect((await store.get(scope, "file-a"))!.generation).toBe( + original.generation, + ); + expect(store.events.filter((event) => event === "publish")).toHaveLength(1); +}); + +test("production embedding requests budget Unicode title and nested headings while preserving source text", async () => { + const { openAICompatible } = await import("../src/embeddings"); + const requests: string[] = []; + const server = Bun.serve({ + hostname: "127.0.0.1", + port: 0, + async fetch(request) { + const body = (await request.json()) as { input: string[] }; + requests.push(...body.input); + if (body.input.some((text) => Buffer.byteLength(text) > 8191)) + return new Response(null, { status: 400 }); + return Response.json({ + data: body.input.map((_, index) => ({ + index, + embedding: [1, 0, 0, 0, 0, 0, 0, 0], + })), + }); + }, + }); + try { + const adapter = openAICompatible({ + endpoint: `http://127.0.0.1:${server.port}/embeddings`, + apiKey: "test", + model: "test", + dimensions: 8, + }); + const store = new MemoryStore(); + const text = + Array.from( + { length: 6 }, + (_, index) => `${"#".repeat(index + 1)} ${"η« ".repeat(512)}\n`, + ).join("") + "ζ–‡".repeat(1500); + const document = await new Pipeline(store, adapter).ingest( + scope, + "file-a", + "user", + "unicode", + { ...input(text), title: "ζ›Έ".repeat(512) }, + signal(), + ); + expect(document.chunkCount).toBeGreaterThan(0); + expect( + requests.every( + (text) => + text.isWellFormed() && + Buffer.byteLength(text) <= adapter.maxInputBytes, + ), + ).toBe(true); + for (const chunk of store.chunks) + expect(chunk.text).toBe(text.slice(chunk.start, chunk.end)); + expect(store.chunks.at(-1)!.section).toHaveLength(6); + } finally { + await server.stop(true); + } +}); From 90c7feef7e3a2364282cf2cfe6405f2bbc230b44 Mon Sep 17 00:00:00 2001 From: Lia Date: Sun, 4 Oct 2026 18:26:51 +0000 Subject: [PATCH 3/3] fix: Qualify the Non-Root RAG Runtime Image --- v2/Dockerfile | 2 +- v2/scripts/qualify.sh | 26 +++++++++- v2/scripts/runtime.ts | 107 ++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 133 insertions(+), 2 deletions(-) create mode 100644 v2/scripts/runtime.ts diff --git a/v2/Dockerfile b/v2/Dockerfile index b8a9f38c..e2c2c992 100644 --- a/v2/Dockerfile +++ b/v2/Dockerfile @@ -9,7 +9,7 @@ FROM oven/bun:1.4.2-debian WORKDIR /app COPY --from=build /app/node_modules ./node_modules COPY --from=build /app/dist ./dist -COPY package.json ./ +COPY --chmod=0444 package.json ./ USER bun ENV RAG_V2_HOST=0.0.0.0 EXPOSE 8001 diff --git a/v2/scripts/qualify.sh b/v2/scripts/qualify.sh index c30307cd..f476331f 100755 --- a/v2/scripts/qualify.sh +++ b/v2/scripts/qualify.sh @@ -4,8 +4,13 @@ cd "$(dirname "$0")/.." run_id="$(date +%s)-$$" ch="ragv2-ch-$run_id" runner="ragv2-bun-$run_id" +app="ragv2-app-$run_id" +image="ragv2-image-$run_id" +manifest_mode=$(stat -c '%a' package.json) cleanup() { - docker rm -f "$runner" "$ch" >/dev/null 2>&1 || true + chmod "$manifest_mode" package.json + docker rm -fv "$app" "$runner" "$ch" >/dev/null 2>&1 || true + docker image rm "$image" >/dev/null 2>&1 || true } trap cleanup EXIT # No published database port or application credentials. @@ -21,6 +26,25 @@ docker run -d --name "$runner" --network "container:$ch" --entrypoint sh \ docker exec "$runner" mkdir /app docker cp . "$runner:/app/v2" docker exec -w /app/v2 "$runner" bun test test/integration.test.ts +# Qualify the shipped non-root image even when checkout files are owner-only. +chmod 600 package.json +docker build --quiet -t "$image" . +chmod "$manifest_mode" package.json +docker exec "$ch" clickhouse-client --query 'CREATE DATABASE rag_v2' >/dev/null +jwks=$(docker exec -w /app/v2 "$runner" bun scripts/runtime.ts prepare) +docker run -d --name "$app" --network "container:$ch" --read-only --cap-drop ALL \ + --security-opt no-new-privileges -e RAG_V2_JWKS_JSON="$jwks" \ + -e RAG_V2_CLICKHOUSE_URL=http://127.0.0.1:8123 \ + -e RAG_V2_EMBEDDING_PROVIDER=test -e RAG_V2_ALLOW_TEST_PROVIDER=true \ + -e RAG_V2_EMBEDDING_DIMENSIONS=64 -e RAG_V2_MODE=coordinator \ + -e RAG_V2_SINGLE_WRITER=true "$image" sh -c 'bun dist/migrate.js && exec bun dist/main.js' >/dev/null +test "$(docker inspect "$app" --format '{{.Config.User}}')" = bun +for attempt in $(seq 1 30); do + if docker exec "$runner" bun --eval 'fetch("http://127.0.0.1:8001/health").then(r=>process.exit(r.ok?0:1)).catch(()=>process.exit(1))' >/dev/null 2>&1; then break; fi + if [ "$attempt" = 30 ]; then docker logs "$app"; exit 1; fi + sleep 1 +done +docker exec -w /app/v2 "$runner" bun scripts/runtime.ts check if [ "${RAG_V2_RUN_BENCH:-false}" = true ]; then docker exec -w /app/v2 "$runner" bun scripts/benchmark.ts fi diff --git a/v2/scripts/runtime.ts b/v2/scripts/runtime.ts new file mode 100644 index 00000000..ae5bbefc --- /dev/null +++ b/v2/scripts/runtime.ts @@ -0,0 +1,107 @@ +import { exportJWK, generateKeyPair, SignJWT } from "jose"; +import { mkdir } from "node:fs/promises"; + +if (process.argv[2] === "prepare") { + const { publicKey, privateKey } = await generateKeyPair("EdDSA"); + const jwk = { + ...(await exportJWK(publicKey)), + kid: "runtime-test", + alg: "EdDSA", + }; + const now = Math.floor(Date.now() / 1000); + const token = await new SignJWT({ + tenant_id: "runtime-test", + actor_kind: "user", + grants: [ + { + namespaceId: "test", + resourceKind: "document", + operations: ["read", "write", "delete"], + }, + ], + }) + .setProtectedHeader({ kid: "runtime-test", alg: "EdDSA", typ: "JWT" }) + .setSubject("test-user") + .setIssuer("librechat") + .setAudience("ragapi") + .setIssuedAt(now) + .setNotBefore(now) + .setExpirationTime(now + 300) + .setJti(crypto.randomUUID()) + .sign(privateKey); + await mkdir(".qualification", { recursive: true }); + await Bun.write(".qualification/runtime-token", token); + // Only public verification material is exported to the service container. + console.log(JSON.stringify({ keys: [jwk] })); +} else if (process.argv[2] === "check") { + const base = "http://127.0.0.1:8001"; + const token = await Bun.file(".qualification/runtime-token").text(); + const path = "/v2/namespaces/test/documents/file"; + const headers = { + Authorization: `Bearer ${token}`, + "Content-Type": "application/json", + "Idempotency-Key": "one", + }; + function check(value: unknown, message: string): asserts value { + if (!value) throw Error(message); + } + const denied = await fetch(base + path); + check(denied.status === 401, "ANONYMOUS_NOT_DENIED"); + const body = { + segments: [{ kind: "page", index: 1, text: "runtime alpha original" }], + original: { + fileId: "file", + revision: "original-revision", + sha256: "a".repeat(64), + filename: "original.pdf", + mediaType: "application/pdf", + }, + }; + const upload = await fetch(base + path, { + method: "PUT", + headers, + body: JSON.stringify(body), + }); + check(upload.status === 200, "RUNTIME_UPLOAD_FAILED"); + const stored = (await upload.json()) as { document: { generation: string } }; + const search = await fetch(base + "/v2/search", { + method: "POST", + headers, + body: JSON.stringify({ + query: "alpha", + namespaces: [{ namespaceId: "test" }], + }), + }); + const result = (await search.json()) as { + hits: Array<{ + original: { revision: string }; + actor: string; + page: number; + }>; + }; + check( + search.ok && + result.hits[0]?.original.revision === body.original.revision && + result.hits[0].actor === "user:test-user" && + result.hits[0].page === 1, + "RUNTIME_SEARCH_FAILED", + ); + const context = await fetch(base + path + "/context", { headers }); + check(context.status === 200, "RUNTIME_CONTEXT_FAILED"); + const deleted = await fetch(base + path, { + method: "DELETE", + headers: { ...headers, "Idempotency-Key": "delete" }, + }); + check(deleted.status === 204, "RUNTIME_DELETE_FAILED"); + const stale = await fetch(base + path, { + method: "PUT", + headers: { ...headers, "Idempotency-Key": "stale" }, + body: JSON.stringify({ ...body, ifMatch: stored.document.generation }), + }); + check(stale.status === 409, "DELETED_GENERATION_RESURRECTED"); + console.log( + "PASS: built non-root Bun service, asymmetric auth, real HTTP ingestion/search/context/delete, original revision/provenance, and stale-edit denial.", + ); +} else { + throw Error("EXPECTED_PREPARE_OR_CHECK"); +}