Sign In

@agentskit/rag

Package Overview
Dependencies
Maintainers
1
Versions
35
Alerts
File Explorer

Advanced tools

Socket logo

Install Socket

Detect and block malicious and high-risk dependencies

Install

@agentskit/rag - npm Package Compare versions

Comparing version
0.4.14
to
0.4.16
+43
dist/chunk-K6CR4AHC.js
// src/chunker.ts
function resolveChunkSize(chunkSize) {
if (!Number.isFinite(chunkSize) || chunkSize <= 0) return Number.POSITIVE_INFINITY;
return Math.floor(chunkSize);
}
function resolveChunkOverlap(chunkOverlap, chunkSize) {
if (!Number.isFinite(chunkOverlap) || chunkOverlap < 0) return 0;
const overlap = Math.floor(chunkOverlap);
if (!Number.isFinite(chunkSize)) return 0;
return Math.min(overlap, Math.max(0, chunkSize - 1));
}
function chunkText(text, options) {
if (!text) return [];
if (options.split) {
return options.split(text).filter((chunk) => chunk.length > 0);
}
const chunkSize = resolveChunkSize(options.chunkSize);
const chunkOverlap = resolveChunkOverlap(options.chunkOverlap, chunkSize);
if (text.length <= chunkSize) return [text];
const chunks = [];
let start = 0;
while (start < text.length) {
let end = Math.min(start + chunkSize, text.length);
if (end < text.length) {
const boundary = text.lastIndexOf(" ", end);
if (boundary > start) {
end = boundary;
}
}
const chunk = text.slice(start, end).trim();
if (chunk.length > 0) {
chunks.push(chunk);
}
if (end >= text.length) break;
const advance = end - start - chunkOverlap;
start += Math.max(advance, 1);
}
return chunks;
}
export { chunkText };
//# sourceMappingURL=chunk-K6CR4AHC.js.map
//# sourceMappingURL=chunk-K6CR4AHC.js.map
{"version":3,"sources":["../src/chunker.ts"],"names":[],"mappings":";AAQA,SAAS,iBAAiB,SAAA,EAA2B;AACnD,EAAA,IAAI,CAAC,OAAO,QAAA,CAAS,SAAS,KAAK,SAAA,IAAa,CAAA,SAAU,MAAA,CAAO,iBAAA;AACjE,EAAA,OAAO,IAAA,CAAK,MAAM,SAAS,CAAA;AAC7B;AAGA,SAAS,mBAAA,CAAoB,cAAsB,SAAA,EAA2B;AAC5E,EAAA,IAAI,CAAC,MAAA,CAAO,QAAA,CAAS,YAAY,CAAA,IAAK,YAAA,GAAe,GAAG,OAAO,CAAA;AAC/D,EAAA,MAAM,OAAA,GAAU,IAAA,CAAK,KAAA,CAAM,YAAY,CAAA;AACvC,EAAA,IAAI,CAAC,MAAA,CAAO,QAAA,CAAS,SAAS,GAAG,OAAO,CAAA;AACxC,EAAA,OAAO,IAAA,CAAK,IAAI,OAAA,EAAS,IAAA,CAAK,IAAI,CAAA,EAAG,SAAA,GAAY,CAAC,CAAC,CAAA;AACrD;AAEO,SAAS,SAAA,CAAU,MAAc,OAAA,EAAiC;AACvE,EAAA,IAAI,CAAC,IAAA,EAAM,OAAO,EAAC;AAEnB,EAAA,IAAI,QAAQ,KAAA,EAAO;AACjB,IAAA,OAAO,OAAA,CAAQ,MAAM,IAAI,CAAA,CAAE,OAAO,CAAA,KAAA,KAAS,KAAA,CAAM,SAAS,CAAC,CAAA;AAAA,EAC7D;AAEA,EAAA,MAAM,SAAA,GAAY,gBAAA,CAAiB,OAAA,CAAQ,SAAS,CAAA;AACpD,EAAA,MAAM,YAAA,GAAe,mBAAA,CAAoB,OAAA,CAAQ,YAAA,EAAc,SAAS,CAAA;AAExE,EAAA,IAAI,IAAA,CAAK,MAAA,IAAU,SAAA,EAAW,OAAO,CAAC,IAAI,CAAA;AAE1C,EAAA,MAAM,SAAmB,EAAC;AAC1B,EAAA,IAAI,KAAA,GAAQ,CAAA;AAEZ,EAAA,OAAO,KAAA,GAAQ,KAAK,MAAA,EAAQ;AAC1B,IAAA,IAAI,MAAM,IAAA,CAAK,GAAA,CAAI,KAAA,GAAQ,SAAA,EAAW,KAAK,MAAM,CAAA;AAEjD,IAAA,IAAI,GAAA,GAAM,KAAK,MAAA,EAAQ;AACrB,MAAA,MAAM,QAAA,GAAW,IAAA,CAAK,WAAA,CAAY,GAAA,EAAK,GAAG,CAAA;AAC1C,MAAA,IAAI,WAAW,KAAA,EAAO;AACpB,QAAA,GAAA,GAAM,QAAA;AAAA,MACR;AAAA,IACF;AAEA,IAAA,MAAM,QAAQ,IAAA,CAAK,KAAA,CAAM,KAAA,EAAO,GAAG,EAAE,IAAA,EAAK;AAC1C,IAAA,IAAI,KAAA,CAAM,SAAS,CAAA,EAAG;AACpB,MAAA,MAAA,CAAO,KAAK,KAAK,CAAA;AAAA,IACnB;AAEA,IAAA,IAAI,GAAA,IAAO,KAAK,MAAA,EAAQ;AAExB,IAAA,MAAM,OAAA,GAAU,MAAM,KAAA,GAAQ,YAAA;AAC9B,IAAA,KAAA,IAAS,IAAA,CAAK,GAAA,CAAI,OAAA,EAAS,CAAC,CAAA;AAAA,EAC9B;AAEA,EAAA,OAAO,MAAA;AACT","file":"chunk-K6CR4AHC.js","sourcesContent":["export interface ChunkOptions {\n chunkSize: number\n chunkOverlap: number\n split?: (text: string) => string[]\n}\n\n// Invalid sizes become +Infinity so the whole text is one chunk and the loop\n// always terminates.\nfunction resolveChunkSize(chunkSize: number): number {\n if (!Number.isFinite(chunkSize) || chunkSize <= 0) return Number.POSITIVE_INFINITY\n return Math.floor(chunkSize)\n}\n\n// Overlap must stay strictly below chunkSize so start always advances.\nfunction resolveChunkOverlap(chunkOverlap: number, chunkSize: number): number {\n if (!Number.isFinite(chunkOverlap) || chunkOverlap < 0) return 0\n const overlap = Math.floor(chunkOverlap)\n if (!Number.isFinite(chunkSize)) return 0\n return Math.min(overlap, Math.max(0, chunkSize - 1))\n}\n\nexport function chunkText(text: string, options: ChunkOptions): string[] {\n if (!text) return []\n\n if (options.split) {\n return options.split(text).filter(chunk => chunk.length > 0)\n }\n\n const chunkSize = resolveChunkSize(options.chunkSize)\n const chunkOverlap = resolveChunkOverlap(options.chunkOverlap, chunkSize)\n\n if (text.length <= chunkSize) return [text]\n\n const chunks: string[] = []\n let start = 0\n\n while (start < text.length) {\n let end = Math.min(start + chunkSize, text.length)\n\n if (end < text.length) {\n const boundary = text.lastIndexOf(' ', end)\n if (boundary > start) {\n end = boundary\n }\n }\n\n const chunk = text.slice(start, end).trim()\n if (chunk.length > 0) {\n chunks.push(chunk)\n }\n\n if (end >= text.length) break\n\n const advance = end - start - chunkOverlap\n start += Math.max(advance, 1)\n }\n\n return chunks\n}\n"]}
+12
-1
'use strict';
// src/chunker.ts
function resolveChunkSize(chunkSize) {
if (!Number.isFinite(chunkSize) || chunkSize <= 0) return Number.POSITIVE_INFINITY;
return Math.floor(chunkSize);
}
function resolveChunkOverlap(chunkOverlap, chunkSize) {
if (!Number.isFinite(chunkOverlap) || chunkOverlap < 0) return 0;
const overlap = Math.floor(chunkOverlap);
if (!Number.isFinite(chunkSize)) return 0;
return Math.min(overlap, Math.max(0, chunkSize - 1));
}
function chunkText(text, options) {

@@ -9,3 +19,4 @@ if (!text) return [];

}
const { chunkSize, chunkOverlap } = options;
const chunkSize = resolveChunkSize(options.chunkSize);
const chunkOverlap = resolveChunkOverlap(options.chunkOverlap, chunkSize);
if (text.length <= chunkSize) return [text];

@@ -12,0 +23,0 @@ const chunks = [];

+1
-1

@@ -1,1 +0,1 @@

{"version":3,"sources":["../src/chunker.ts"],"names":[],"mappings":";;;AAMO,SAAS,SAAA,CAAU,MAAc,OAAA,EAAiC;AACvE,EAAA,IAAI,CAAC,IAAA,EAAM,OAAO,EAAC;AAEnB,EAAA,IAAI,QAAQ,KAAA,EAAO;AACjB,IAAA,OAAO,OAAA,CAAQ,MAAM,IAAI,CAAA,CAAE,OAAO,CAAA,KAAA,KAAS,KAAA,CAAM,SAAS,CAAC,CAAA;AAAA,EAC7D;AAEA,EAAA,MAAM,EAAE,SAAA,EAAW,YAAA,EAAa,GAAI,OAAA;AAEpC,EAAA,IAAI,IAAA,CAAK,MAAA,IAAU,SAAA,EAAW,OAAO,CAAC,IAAI,CAAA;AAE1C,EAAA,MAAM,SAAmB,EAAC;AAC1B,EAAA,IAAI,KAAA,GAAQ,CAAA;AAEZ,EAAA,OAAO,KAAA,GAAQ,KAAK,MAAA,EAAQ;AAC1B,IAAA,IAAI,MAAM,IAAA,CAAK,GAAA,CAAI,KAAA,GAAQ,SAAA,EAAW,KAAK,MAAM,CAAA;AAEjD,IAAA,IAAI,GAAA,GAAM,KAAK,MAAA,EAAQ;AACrB,MAAA,MAAM,QAAA,GAAW,IAAA,CAAK,WAAA,CAAY,GAAA,EAAK,GAAG,CAAA;AAC1C,MAAA,IAAI,WAAW,KAAA,EAAO;AACpB,QAAA,GAAA,GAAM,QAAA;AAAA,MACR;AAAA,IACF;AAEA,IAAA,MAAM,QAAQ,IAAA,CAAK,KAAA,CAAM,KAAA,EAAO,GAAG,EAAE,IAAA,EAAK;AAC1C,IAAA,IAAI,KAAA,CAAM,SAAS,CAAA,EAAG;AACpB,MAAA,MAAA,CAAO,KAAK,KAAK,CAAA;AAAA,IACnB;AAEA,IAAA,IAAI,GAAA,IAAO,KAAK,MAAA,EAAQ;AAExB,IAAA,MAAM,OAAA,GAAU,MAAM,KAAA,GAAQ,YAAA;AAC9B,IAAA,KAAA,IAAS,IAAA,CAAK,GAAA,CAAI,OAAA,EAAS,CAAC,CAAA;AAAA,EAC9B;AAEA,EAAA,OAAO,MAAA;AACT","file":"chunker.cjs","sourcesContent":["export interface ChunkOptions {\n chunkSize: number\n chunkOverlap: number\n split?: (text: string) => string[]\n}\n\nexport function chunkText(text: string, options: ChunkOptions): string[] {\n if (!text) return []\n\n if (options.split) {\n return options.split(text).filter(chunk => chunk.length > 0)\n }\n\n const { chunkSize, chunkOverlap } = options\n\n if (text.length <= chunkSize) return [text]\n\n const chunks: string[] = []\n let start = 0\n\n while (start < text.length) {\n let end = Math.min(start + chunkSize, text.length)\n\n if (end < text.length) {\n const boundary = text.lastIndexOf(' ', end)\n if (boundary > start) {\n end = boundary\n }\n }\n\n const chunk = text.slice(start, end).trim()\n if (chunk.length > 0) {\n chunks.push(chunk)\n }\n\n if (end >= text.length) break\n\n const advance = end - start - chunkOverlap\n start += Math.max(advance, 1)\n }\n\n return chunks\n}\n"]}
{"version":3,"sources":["../src/chunker.ts"],"names":[],"mappings":";;;AAQA,SAAS,iBAAiB,SAAA,EAA2B;AACnD,EAAA,IAAI,CAAC,OAAO,QAAA,CAAS,SAAS,KAAK,SAAA,IAAa,CAAA,SAAU,MAAA,CAAO,iBAAA;AACjE,EAAA,OAAO,IAAA,CAAK,MAAM,SAAS,CAAA;AAC7B;AAGA,SAAS,mBAAA,CAAoB,cAAsB,SAAA,EAA2B;AAC5E,EAAA,IAAI,CAAC,MAAA,CAAO,QAAA,CAAS,YAAY,CAAA,IAAK,YAAA,GAAe,GAAG,OAAO,CAAA;AAC/D,EAAA,MAAM,OAAA,GAAU,IAAA,CAAK,KAAA,CAAM,YAAY,CAAA;AACvC,EAAA,IAAI,CAAC,MAAA,CAAO,QAAA,CAAS,SAAS,GAAG,OAAO,CAAA;AACxC,EAAA,OAAO,IAAA,CAAK,IAAI,OAAA,EAAS,IAAA,CAAK,IAAI,CAAA,EAAG,SAAA,GAAY,CAAC,CAAC,CAAA;AACrD;AAEO,SAAS,SAAA,CAAU,MAAc,OAAA,EAAiC;AACvE,EAAA,IAAI,CAAC,IAAA,EAAM,OAAO,EAAC;AAEnB,EAAA,IAAI,QAAQ,KAAA,EAAO;AACjB,IAAA,OAAO,OAAA,CAAQ,MAAM,IAAI,CAAA,CAAE,OAAO,CAAA,KAAA,KAAS,KAAA,CAAM,SAAS,CAAC,CAAA;AAAA,EAC7D;AAEA,EAAA,MAAM,SAAA,GAAY,gBAAA,CAAiB,OAAA,CAAQ,SAAS,CAAA;AACpD,EAAA,MAAM,YAAA,GAAe,mBAAA,CAAoB,OAAA,CAAQ,YAAA,EAAc,SAAS,CAAA;AAExE,EAAA,IAAI,IAAA,CAAK,MAAA,IAAU,SAAA,EAAW,OAAO,CAAC,IAAI,CAAA;AAE1C,EAAA,MAAM,SAAmB,EAAC;AAC1B,EAAA,IAAI,KAAA,GAAQ,CAAA;AAEZ,EAAA,OAAO,KAAA,GAAQ,KAAK,MAAA,EAAQ;AAC1B,IAAA,IAAI,MAAM,IAAA,CAAK,GAAA,CAAI,KAAA,GAAQ,SAAA,EAAW,KAAK,MAAM,CAAA;AAEjD,IAAA,IAAI,GAAA,GAAM,KAAK,MAAA,EAAQ;AACrB,MAAA,MAAM,QAAA,GAAW,IAAA,CAAK,WAAA,CAAY,GAAA,EAAK,GAAG,CAAA;AAC1C,MAAA,IAAI,WAAW,KAAA,EAAO;AACpB,QAAA,GAAA,GAAM,QAAA;AAAA,MACR;AAAA,IACF;AAEA,IAAA,MAAM,QAAQ,IAAA,CAAK,KAAA,CAAM,KAAA,EAAO,GAAG,EAAE,IAAA,EAAK;AAC1C,IAAA,IAAI,KAAA,CAAM,SAAS,CAAA,EAAG;AACpB,MAAA,MAAA,CAAO,KAAK,KAAK,CAAA;AAAA,IACnB;AAEA,IAAA,IAAI,GAAA,IAAO,KAAK,MAAA,EAAQ;AAExB,IAAA,MAAM,OAAA,GAAU,MAAM,KAAA,GAAQ,YAAA;AAC9B,IAAA,KAAA,IAAS,IAAA,CAAK,GAAA,CAAI,OAAA,EAAS,CAAC,CAAA;AAAA,EAC9B;AAEA,EAAA,OAAO,MAAA;AACT","file":"chunker.cjs","sourcesContent":["export interface ChunkOptions {\n chunkSize: number\n chunkOverlap: number\n split?: (text: string) => string[]\n}\n\n// Invalid sizes become +Infinity so the whole text is one chunk and the loop\n// always terminates.\nfunction resolveChunkSize(chunkSize: number): number {\n if (!Number.isFinite(chunkSize) || chunkSize <= 0) return Number.POSITIVE_INFINITY\n return Math.floor(chunkSize)\n}\n\n// Overlap must stay strictly below chunkSize so start always advances.\nfunction resolveChunkOverlap(chunkOverlap: number, chunkSize: number): number {\n if (!Number.isFinite(chunkOverlap) || chunkOverlap < 0) return 0\n const overlap = Math.floor(chunkOverlap)\n if (!Number.isFinite(chunkSize)) return 0\n return Math.min(overlap, Math.max(0, chunkSize - 1))\n}\n\nexport function chunkText(text: string, options: ChunkOptions): string[] {\n if (!text) return []\n\n if (options.split) {\n return options.split(text).filter(chunk => chunk.length > 0)\n }\n\n const chunkSize = resolveChunkSize(options.chunkSize)\n const chunkOverlap = resolveChunkOverlap(options.chunkOverlap, chunkSize)\n\n if (text.length <= chunkSize) return [text]\n\n const chunks: string[] = []\n let start = 0\n\n while (start < text.length) {\n let end = Math.min(start + chunkSize, text.length)\n\n if (end < text.length) {\n const boundary = text.lastIndexOf(' ', end)\n if (boundary > start) {\n end = boundary\n }\n }\n\n const chunk = text.slice(start, end).trim()\n if (chunk.length > 0) {\n chunks.push(chunk)\n }\n\n if (end >= text.length) break\n\n const advance = end - start - chunkOverlap\n start += Math.max(advance, 1)\n }\n\n return chunks\n}\n"]}

@@ -1,3 +0,3 @@

export { chunkText } from './chunk-UNFVK5RA.js';
export { chunkText } from './chunk-K6CR4AHC.js';
//# sourceMappingURL=chunker.js.map
//# sourceMappingURL=chunker.js.map

@@ -8,2 +8,12 @@ 'use strict';

// src/chunker.ts
function resolveChunkSize(chunkSize) {
if (!Number.isFinite(chunkSize) || chunkSize <= 0) return Number.POSITIVE_INFINITY;
return Math.floor(chunkSize);
}
function resolveChunkOverlap(chunkOverlap, chunkSize) {
if (!Number.isFinite(chunkOverlap) || chunkOverlap < 0) return 0;
const overlap = Math.floor(chunkOverlap);
if (!Number.isFinite(chunkSize)) return 0;
return Math.min(overlap, Math.max(0, chunkSize - 1));
}
function chunkText(text, options) {

@@ -14,3 +24,4 @@ if (!text) return [];

}
const { chunkSize, chunkOverlap } = options;
const chunkSize = resolveChunkSize(options.chunkSize);
const chunkOverlap = resolveChunkOverlap(options.chunkOverlap, chunkSize);
if (text.length <= chunkSize) return [text];

@@ -39,2 +50,31 @@ const chunks = [];

// src/rag.ts
function resolvePositiveInt(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return Math.max(1, Math.floor(value));
}
function resolveThreshold(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return value;
}
function resolveChunkSize2(value, fallback) {
if (value === void 0 || !Number.isFinite(value) || value <= 0) return fallback;
return Math.floor(value);
}
function resolveChunkOverlap2(value, fallback, chunkSize) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return fallback;
return Math.min(Math.floor(value), Math.max(0, chunkSize - 1));
}
function enforceScoreOrder(docs) {
if (docs.length === 0) return docs;
const anyScore = docs.some((d) => d.score !== void 0);
if (!anyScore) return docs;
for (const d of docs) {
if (typeof d.score !== "number" || !Number.isFinite(d.score)) {
throw new TypeError(
"createRAG search: every result must have a finite numeric score when scores are present"
);
}
}
return [...docs].sort((a, b) => b.score - a.score);
}
function createRAG(config) {

@@ -44,13 +84,17 @@ const {

store,
chunkSize = 512,
chunkOverlap = 50,
split,
topK = 5,
threshold = 0
split
} = config;
const chunkSize = resolveChunkSize2(config.chunkSize, 512);
const chunkOverlap = resolveChunkOverlap2(config.chunkOverlap, 50, chunkSize);
const defaultTopK = resolvePositiveInt(config.topK, 5);
const defaultThreshold = resolveThreshold(config.threshold, 0);
async function ingest(documents) {
const vectorDocs = [];
for (const doc of documents) {
const content = doc.content;
if (!content) continue;
const docId = doc.id ?? core.generateId("doc");
const chunks = chunkText(doc.content, { chunkSize, chunkOverlap, split });
const source = doc.source;
const metadata = doc.metadata ? { ...doc.metadata } : void 0;
const chunks = chunkText(content, { chunkSize, chunkOverlap, split });
for (let i = 0; i < chunks.length; i++) {

@@ -64,4 +108,4 @@ const chunkId = `${docId}_chunk_${i}`;

metadata: {
...doc.metadata,
source: doc.source,
...metadata,
source,
documentId: docId,

@@ -78,7 +122,7 @@ chunkIndex: i

async function search(query, options) {
const topK = options?.topK !== void 0 ? resolvePositiveInt(options.topK, defaultTopK) : defaultTopK;
const threshold = options?.threshold !== void 0 ? resolveThreshold(options.threshold, defaultThreshold) : defaultThreshold;
const queryEmbedding = await embed(query);
return store.search(queryEmbedding, {
topK: options?.topK ?? topK,
threshold: options?.threshold ?? threshold
});
const results = await store.search(queryEmbedding, { topK, threshold });
return enforceScoreOrder(results);
}

@@ -105,8 +149,100 @@ async function retrieve(request) {

// src/loaders.ts
// src/loaders/shared.ts
function resolveMaxFiles(maxFiles, fallback = 100) {
if (maxFiles === void 0) return fallback;
if (!Number.isFinite(maxFiles)) return 0;
return Math.max(0, Math.floor(maxFiles));
}
function loadFailed(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_LOAD_FAILED,
message,
cause
});
}
function isAbortLike(err) {
if (err == null || typeof err !== "object") return false;
const name = err.name;
if (name === "AbortError") return true;
const cause = err.cause;
return cause !== void 0 && isAbortLike(cause);
}
function ensureNotAborted(signal, label) {
if (signal?.aborted) {
throw loadFailed(`${label}: aborted`, signal.reason);
}
}
function rethrowIfAbort(err, signal, label) {
ensureNotAborted(signal, label);
if (isAbortLike(err)) {
if (err instanceof RagError) throw err;
throw loadFailed(`${label}: aborted`, err);
}
}
function finishTreeLoad(label, attempted, loaded, docs) {
if (attempted > 0 && loaded === 0) {
throw loadFailed(`${label}: all eligible downloads failed`);
}
return docs;
}
async function doFetch(fetchImpl, url, init, label) {
try {
return await fetchImpl(url, init);
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) {
throw loadFailed(`${label}: aborted`, cause);
}
throw loadFailed(`${label}: network error for ${url}`, cause);
}
}
async function readResponseText(response, label) {
try {
return await response.text();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read response body`, cause);
}
}
async function readResponseJson(response, label) {
try {
return await response.json();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to parse response body`, cause);
}
}
async function readResponseArrayBuffer(response, label) {
try {
return await response.arrayBuffer();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read response body`, cause);
}
}
async function readS3Body(body, label) {
if (body == null || typeof body.transformToString !== "function") {
throw loadFailed(`${label}: missing or invalid object body`);
}
try {
return await body.transformToString();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read object body`, cause);
}
}
function encodePathSegments(path) {
return path.split("/").map((segment) => encodeURIComponent(segment)).join("/");
}
// src/loaders/documents.ts
async function loadUrl(url, options = {}) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await fetchImpl(url, { headers: options.headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadUrl ${response.status}: ${url}` });
const content = await response.text();
const response = await doFetch(fetchImpl, url, { headers: options.headers, signal: options.signal }, "loadUrl");
if (!response.ok) throw loadFailed(`loadUrl ${response.status}: ${url}`);
const content = await readResponseText(response, "loadUrl");
return [{ content, source: url, metadata: { url } }];

@@ -117,10 +253,11 @@ }

const ref = options.ref ?? "HEAD";
const url = `https://raw.githubusercontent.com/${owner}/${repo}/${ref}/${path}`;
const encodedPath = encodePathSegments(path);
const url = `https://raw.githubusercontent.com/${owner}/${repo}/${ref}/${encodedPath}`;
const headers = {};
if (options.token) headers.authorization = `Bearer ${options.token}`;
const response = await fetchImpl(url, { headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGitHubFile ${response.status}: ${url}` });
const response = await doFetch(fetchImpl, url, { headers, signal: options.signal }, "loadGitHubFile");
if (!response.ok) throw loadFailed(`loadGitHubFile ${response.status}: ${url}`);
return [
{
content: await response.text(),
content: await readResponseText(response, "loadGitHubFile"),
source: url,

@@ -134,29 +271,61 @@ metadata: { owner, repo, path, ref }

const ref = options.ref ?? "HEAD";
const url = `https://api.github.com/repos/${owner}/${repo}/git/trees/${ref}?recursive=1`;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const url = `https://api.github.com/repos/${owner}/${repo}/git/trees/${encodeURIComponent(ref)}?recursive=1`;
const headers = { accept: "application/vnd.github+json" };
if (options.token) headers.authorization = `Bearer ${options.token}`;
const response = await fetchImpl(url, { headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGitHubTree ${response.status}: ${url}` });
const tree = await response.json();
const files = (tree.tree ?? []).filter((t) => t.type === "blob").filter((t) => !options.filter || options.filter(t.path)).slice(0, options.maxFiles ?? 100);
const response = await doFetch(fetchImpl, url, { headers, signal: options.signal }, "loadGitHubTree");
if (!response.ok) throw loadFailed(`loadGitHubTree ${response.status}: ${url}`);
const tree = await readResponseJson(
response,
"loadGitHubTree"
);
const files = (tree.tree ?? []).filter((t) => t.type === "blob").filter((t) => !options.filter || options.filter(t.path)).slice(0, maxFiles);
const docs = [];
let attempted = 0;
let loaded = 0;
for (const file of files) {
const items = await loadGitHubFile(owner, repo, file.path, options);
docs.push(...items);
ensureNotAborted(options.signal, "loadGitHubTree");
attempted++;
try {
const items = await loadGitHubFile(owner, repo, file.path, options);
docs.push(...items);
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadGitHubTree");
}
}
return docs;
return finishTreeLoad("loadGitHubTree", attempted, loaded, docs);
}
async function loadNotionPage(pageId, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const url = `https://api.notion.com/v1/blocks/${pageId}/children?page_size=100`;
const response = await fetchImpl(url, {
headers: {
authorization: `Bearer ${options.token}`,
"notion-version": options.version ?? "2022-06-28"
const HEADING_PREFIX = { heading_1: "# ", heading_2: "## ", heading_3: "### " };
const blocks = [];
let cursor;
const seenCursors = /* @__PURE__ */ new Set();
while (true) {
ensureNotAborted(options.signal, "loadNotionPage");
let url = `https://api.notion.com/v1/blocks/${pageId}/children?page_size=100`;
if (cursor) url += `&start_cursor=${encodeURIComponent(cursor)}`;
const response = await doFetch(fetchImpl, url, {
headers: {
authorization: `Bearer ${options.token}`,
"notion-version": options.version ?? "2022-06-28"
},
signal: options.signal
}, "loadNotionPage");
if (!response.ok) throw loadFailed(`loadNotionPage ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadNotionPage");
for (const block of data.results ?? []) {
blocks.push(block);
}
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadNotionPage ${response.status}: ${url}` });
const data = await response.json();
const HEADING_PREFIX = { heading_1: "# ", heading_2: "## ", heading_3: "### " };
const text = (data.results ?? []).map((block) => {
if (!data.has_more) break;
const next = data.next_cursor;
if (typeof next !== "string" || next.length === 0 || seenCursors.has(next)) {
throw loadFailed("loadNotionPage: incomplete pagination (has_more without a new cursor)");
}
seenCursors.add(next);
cursor = next;
}
const text = blocks.map((block) => {
const part = block.paragraph?.rich_text ?? block.heading_1?.rich_text ?? block.heading_2?.rich_text ?? block.heading_3?.rich_text;

@@ -173,7 +342,11 @@ if (!part) return "";

const authHeader = options.authorization ?? (options.token ? `Basic ${options.token}` : void 0);
const response = await fetchImpl(url, {
headers: authHeader ? { authorization: authHeader } : {}
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadConfluencePage ${response.status}: ${url}` });
const data = await response.json();
const response = await doFetch(fetchImpl, url, {
headers: authHeader ? { authorization: authHeader } : {},
signal: options.signal
}, "loadConfluencePage");
if (!response.ok) throw loadFailed(`loadConfluencePage ${response.status}: ${url}`);
const data = await readResponseJson(
response,
"loadConfluencePage"
);
const content = data.body?.storage?.value ?? "";

@@ -185,9 +358,20 @@ return [{ content, source: `${options.baseUrl}/pages/${pageId}`, metadata: { pageId, title: data.title } }];

const url = `https://www.googleapis.com/drive/v3/files/${fileId}/export?mimeType=text/plain`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${options.accessToken}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGoogleDriveFile ${response.status}: ${url}` });
const content = await response.text();
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${options.accessToken}` },
signal: options.signal
}, "loadGoogleDriveFile");
if (!response.ok) throw loadFailed(`loadGoogleDriveFile ${response.status}: ${url}`);
const content = await readResponseText(response, "loadGoogleDriveFile");
return [{ content, source: `gdrive://${fileId}`, metadata: { fileId } }];
}
async function loadPdf(url, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await doFetch(fetchImpl, url, { signal: options.signal }, "loadPdf");
if (!response.ok) throw loadFailed(`loadPdf ${response.status}: ${url}`);
const buf = new Uint8Array(await readResponseArrayBuffer(response, "loadPdf"));
const { text, pages } = await options.parsePdf(buf);
return [{ content: text, source: url, metadata: { url, pages } }];
}
// src/loaders/s3.ts
async function loadS3(options) {

@@ -201,67 +385,108 @@ if (!options.commands) {

}
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const { ListObjectsV2Command, GetObjectCommand } = options.commands;
const docs = [];
let attempted = 0;
let loaded = 0;
let continuationToken;
const maxFiles = options.maxFiles ?? 100;
outer: while (true) {
const list = await options.client.send(new ListObjectsV2Command({
Bucket: options.bucket,
Prefix: options.prefix,
ContinuationToken: continuationToken
}));
ensureNotAborted(options.signal, "loadS3");
let list;
try {
list = await options.client.send(new ListObjectsV2Command({
Bucket: options.bucket,
Prefix: options.prefix,
ContinuationToken: continuationToken
}));
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed("loadS3: aborted", cause);
throw loadFailed(`loadS3: list failed for s3://${options.bucket}`, cause);
}
for (const obj of list.Contents ?? []) {
ensureNotAborted(options.signal, "loadS3");
const key = obj.Key;
if (!key) continue;
if (options.filter && !options.filter(key)) continue;
const get = await options.client.send(new GetObjectCommand({
Bucket: options.bucket,
Key: key
}));
const content = await get.Body?.transformToString() ?? "";
docs.push({
content,
source: `s3://${options.bucket}/${key}`,
metadata: { bucket: options.bucket, key }
});
attempted++;
try {
const get = await options.client.send(new GetObjectCommand({
Bucket: options.bucket,
Key: key
}));
const content = await readS3Body(get.Body, "loadS3");
docs.push({
content,
source: `s3://${options.bucket}/${key}`,
metadata: { bucket: options.bucket, key }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadS3");
}
if (docs.length >= maxFiles) break outer;
}
if (!list.IsTruncated) break;
continuationToken = list.NextContinuationToken;
const next = list.NextContinuationToken;
if (!next || next === continuationToken) {
throw loadFailed("loadS3: incomplete pagination (truncated without a new continuation token)");
}
continuationToken = next;
}
return docs;
return finishTreeLoad("loadS3", attempted, loaded, docs);
}
// src/loaders/cloud.ts
async function loadGcs(options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const docs = [];
const maxFiles = options.maxFiles ?? 100;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const getToken = async () => typeof options.accessToken === "string" ? options.accessToken : await options.accessToken();
let attempted = 0;
let loaded = 0;
let pageToken;
outer: while (true) {
ensureNotAborted(options.signal, "loadGcs");
const params = new URLSearchParams();
if (options.prefix) params.set("prefix", options.prefix);
if (pageToken) params.set("pageToken", pageToken);
const url = `https://storage.googleapis.com/storage/v1/b/${options.bucket}/o?${params.toString()}`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGcs ${response.status}: ${url}` });
const data = await response.json();
const url = `https://storage.googleapis.com/storage/v1/b/${encodeURIComponent(options.bucket)}/o?${params.toString()}`;
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadGcs");
if (!response.ok) throw loadFailed(`loadGcs ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadGcs");
for (const item of data.items ?? []) {
ensureNotAborted(options.signal, "loadGcs");
if (options.filter && !options.filter(item.name)) continue;
const objUrl = `https://storage.googleapis.com/storage/v1/b/${options.bucket}/o/${encodeURIComponent(item.name)}?alt=media`;
const objResponse = await fetchImpl(objUrl, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!objResponse.ok) continue;
docs.push({
content: await objResponse.text(),
source: `gs://${options.bucket}/${item.name}`,
metadata: { bucket: options.bucket, name: item.name }
});
const objUrl = `https://storage.googleapis.com/storage/v1/b/${encodeURIComponent(options.bucket)}/o/${encodeURIComponent(item.name)}?alt=media`;
attempted++;
try {
const objResponse = await doFetch(fetchImpl, objUrl, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadGcs");
if (!objResponse.ok) continue;
docs.push({
content: await readResponseText(objResponse, "loadGcs"),
source: `gs://${options.bucket}/${item.name}`,
metadata: { bucket: options.bucket, name: item.name }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadGcs");
}
if (docs.length >= maxFiles) break outer;
}
if (!data.nextPageToken) break;
pageToken = data.nextPageToken;
const next = data.nextPageToken;
if (!next) break;
if (next === pageToken) {
throw loadFailed("loadGcs: incomplete pagination (no new page token)");
}
pageToken = next;
}
return docs;
return finishTreeLoad("loadGcs", attempted, loaded, docs);
}

@@ -271,12 +496,22 @@ async function loadDropbox(options) {

const docs = [];
const maxFiles = options.maxFiles ?? 100;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const headers = { authorization: `Bearer ${options.accessToken}`, "content-type": "application/json" };
let attempted = 0;
let loaded = 0;
let cursor;
outer: while (true) {
ensureNotAborted(options.signal, "loadDropbox");
const url = cursor ? "https://api.dropboxapi.com/2/files/list_folder/continue" : "https://api.dropboxapi.com/2/files/list_folder";
const body = cursor ? { cursor } : { path: options.path ?? "", recursive: true };
const response = await fetchImpl(url, { method: "POST", headers, body: JSON.stringify(body) });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadDropbox ${response.status}: ${url}` });
const data = await response.json();
const response = await doFetch(fetchImpl, url, {
method: "POST",
headers,
body: JSON.stringify(body),
signal: options.signal
}, "loadDropbox");
if (!response.ok) throw loadFailed(`loadDropbox ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadDropbox");
for (const entry of data.entries ?? []) {
ensureNotAborted(options.signal, "loadDropbox");
if (entry[".tag"] !== "file") continue;

@@ -286,21 +521,32 @@ const path = entry.path_display ?? entry.path_lower;

if (options.filter && !options.filter(path)) continue;
const downloadResponse = await fetchImpl("https://content.dropboxapi.com/2/files/download", {
method: "POST",
headers: {
authorization: `Bearer ${options.accessToken}`,
"Dropbox-API-Arg": JSON.stringify({ path })
}
});
if (!downloadResponse.ok) continue;
docs.push({
content: await downloadResponse.text(),
source: `dropbox:${path}`,
metadata: { path }
});
attempted++;
try {
const downloadResponse = await doFetch(fetchImpl, "https://content.dropboxapi.com/2/files/download", {
method: "POST",
headers: {
authorization: `Bearer ${options.accessToken}`,
"Dropbox-API-Arg": JSON.stringify({ path })
},
signal: options.signal
}, "loadDropbox");
if (!downloadResponse.ok) continue;
docs.push({
content: await readResponseText(downloadResponse, "loadDropbox"),
source: `dropbox:${path}`,
metadata: { path }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadDropbox");
}
if (docs.length >= maxFiles) break outer;
}
if (!data.has_more) break;
cursor = data.cursor;
const next = data.cursor;
if (!next || next === cursor) {
throw loadFailed("loadDropbox: incomplete pagination (has_more without a new cursor)");
}
cursor = next;
}
return docs;
return finishTreeLoad("loadDropbox", attempted, loaded, docs);
}

@@ -310,45 +556,77 @@ async function loadOneDrive(options) {

const docs = [];
const maxFiles = options.maxFiles ?? 100;
const driveBase = options.driveId ? `https://graph.microsoft.com/v1.0/drives/${options.driveId}` : `https://graph.microsoft.com/v1.0/me/drive`;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const driveBase = options.driveId ? `https://graph.microsoft.com/v1.0/drives/${encodeURIComponent(options.driveId)}` : `https://graph.microsoft.com/v1.0/me/drive`;
const folder = options.folderItemId ? `items/${options.folderItemId}` : "root";
const getToken = async () => typeof options.accessToken === "string" ? options.accessToken : await options.accessToken();
const visitedFolders = /* @__PURE__ */ new Set();
let attempted = 0;
let loaded = 0;
async function walk(prefix) {
const url = `${driveBase}/${prefix}/children`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadOneDrive ${response.status}: ${url}` });
const data = await response.json();
for (const item of data.value ?? []) {
if (docs.length >= maxFiles) return;
if (visitedFolders.has(prefix)) return;
visitedFolders.add(prefix);
let url = `${driveBase}/${prefix}/children`;
const seenLinks = /* @__PURE__ */ new Set();
while (url) {
ensureNotAborted(options.signal, "loadOneDrive");
if (docs.length >= maxFiles) return;
if (item.folder) {
await walk(`items/${item.id}`);
if (seenLinks.has(url)) {
throw loadFailed("loadOneDrive: incomplete pagination (repeated @odata.nextLink)");
}
seenLinks.add(url);
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadOneDrive");
if (!response.ok) throw loadFailed(`loadOneDrive ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadOneDrive");
for (const item of data.value ?? []) {
ensureNotAborted(options.signal, "loadOneDrive");
if (docs.length >= maxFiles) return;
if (item.folder) {
await walk(`items/${item.id}`);
continue;
}
if (!item.file) continue;
if (options.filter && !options.filter(item.name)) continue;
const downloadUrl = item["@microsoft.graph.downloadUrl"];
if (!downloadUrl) continue;
attempted++;
try {
const fileResponse = await doFetch(fetchImpl, downloadUrl, { signal: options.signal }, "loadOneDrive");
if (!fileResponse.ok) continue;
docs.push({
content: await readResponseText(fileResponse, "loadOneDrive"),
source: `onedrive:${item.id}`,
metadata: { id: item.id, name: item.name, mimeType: item.file.mimeType }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadOneDrive");
}
}
const next = data["@odata.nextLink"];
if (!next) {
url = void 0;
continue;
}
if (!item.file) continue;
if (options.filter && !options.filter(item.name)) continue;
const downloadUrl = item["@microsoft.graph.downloadUrl"];
if (!downloadUrl) continue;
const fileResponse = await fetchImpl(downloadUrl);
if (!fileResponse.ok) continue;
docs.push({
content: await fileResponse.text(),
source: `onedrive:${item.id}`,
metadata: { id: item.id, name: item.name, mimeType: item.file.mimeType }
});
if (seenLinks.has(next) || next === url) {
throw loadFailed("loadOneDrive: incomplete pagination (repeated @odata.nextLink)");
}
url = next;
}
}
await walk(folder);
return docs;
return finishTreeLoad("loadOneDrive", attempted, loaded, docs);
}
async function loadPdf(url, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await fetchImpl(url);
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadPdf ${response.status}: ${url}` });
const buf = new Uint8Array(await response.arrayBuffer());
const { text, pages } = await options.parsePdf(buf);
return [{ content: text, source: url, metadata: { url, pages } }];
}
// src/rerankers/voyage.ts
function rerankFailed(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function voyageReranker(options) {

@@ -359,23 +637,52 @@ const fetchImpl = options.fetch ?? globalThis.fetch;

if (documents.length === 0) return documents;
const response = await fetchImpl("https://api.voyageai.com/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
})
});
let response;
try {
response = await fetchImpl("https://api.voyageai.com/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
}),
signal: options.signal
});
} catch (cause) {
throw rerankFailed("voyage rerank: network error", cause);
}
if (!response.ok) {
const text = await response.text();
throw new RagError({ code: RagErrorCodes.AK_RAG_RERANK_FAILED, message: `voyage rerank: ${response.status} ${text.slice(0, 200)}` });
let text = "";
try {
text = await response.text();
} catch (cause) {
throw rerankFailed(`voyage rerank: ${response.status} (failed to read error body)`, cause);
}
throw rerankFailed(`voyage rerank: ${response.status} ${text.slice(0, 200)}`);
}
const data = await response.json();
const ranked = (data.data ?? []).sort((a, b) => b.relevance_score - a.relevance_score).map((r) => {
let data;
try {
data = await response.json();
} catch (cause) {
throw rerankFailed("voyage rerank: invalid JSON response", cause);
}
if (!Array.isArray(data.data)) {
throw rerankFailed("voyage rerank: data must be an array");
}
const ranked = [];
for (let i = 0; i < data.data.length; i++) {
const r = data.data[i];
if (r == null || typeof r !== "object" || !Number.isInteger(r.index) || r.index < 0 || r.index >= documents.length || typeof r.relevance_score !== "number" || !Number.isFinite(r.relevance_score)) {
throw rerankFailed(`voyage rerank: malformed result at index ${i}`);
}
const doc = documents[r.index];
return { ...doc, score: r.relevance_score };
});
ranked.push({
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0,
score: r.relevance_score
});
}
ranked.sort((a, b) => b.score - a.score);
return ranked;

@@ -386,2 +693,9 @@ };

// src/rerankers/jina.ts
function rerankFailed2(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function jinaReranker(options) {

@@ -392,23 +706,52 @@ const fetchImpl = options.fetch ?? globalThis.fetch;

if (documents.length === 0) return documents;
const response = await fetchImpl("https://api.jina.ai/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
})
});
let response;
try {
response = await fetchImpl("https://api.jina.ai/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
}),
signal: options.signal
});
} catch (cause) {
throw rerankFailed2("jina rerank: network error", cause);
}
if (!response.ok) {
const text = await response.text();
throw new RagError({ code: RagErrorCodes.AK_RAG_RERANK_FAILED, message: `jina rerank: ${response.status} ${text.slice(0, 200)}` });
let text = "";
try {
text = await response.text();
} catch (cause) {
throw rerankFailed2(`jina rerank: ${response.status} (failed to read error body)`, cause);
}
throw rerankFailed2(`jina rerank: ${response.status} ${text.slice(0, 200)}`);
}
const data = await response.json();
const ranked = (data.results ?? []).sort((a, b) => b.relevance_score - a.relevance_score).map((r) => {
let data;
try {
data = await response.json();
} catch (cause) {
throw rerankFailed2("jina rerank: invalid JSON response", cause);
}
if (!Array.isArray(data.results)) {
throw rerankFailed2("jina rerank: results must be an array");
}
const ranked = [];
for (let i = 0; i < data.results.length; i++) {
const r = data.results[i];
if (r == null || typeof r !== "object" || !Number.isInteger(r.index) || r.index < 0 || r.index >= documents.length || typeof r.relevance_score !== "number" || !Number.isFinite(r.relevance_score)) {
throw rerankFailed2(`jina rerank: malformed result at index ${i}`);
}
const doc = documents[r.index];
return { ...doc, score: r.relevance_score };
});
ranked.push({
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0,
score: r.relevance_score
});
}
ranked.sort((a, b) => b.score - a.score);
return ranked;

@@ -419,11 +762,85 @@ };

// src/rerank.ts
function resolvePositiveInt2(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return Math.max(1, Math.floor(value));
}
function resolveWeight(value, fallback) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return fallback;
return value;
}
function resolveRelativeWeights(vectorWeight, bm25Weight) {
const v = resolveWeight(vectorWeight, 0.6);
const b = resolveWeight(bm25Weight, 0.4);
if (v === 0 && b === 0) return { vectorWeight: 0.5, bm25Weight: 0.5 };
const m = Math.max(v, b);
const vN = v / m;
const bN = b / m;
const sum = vN + bN;
return { vectorWeight: vN / sum, bm25Weight: bN / sum };
}
function cloneDoc(doc) {
return {
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0
};
}
function rerankFailed3(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function enforceScoreOrder2(docs) {
if (docs.length === 0) return docs;
const anyScore = docs.some((d) => d.score !== void 0);
if (!anyScore) return docs.map(cloneDoc);
for (const d of docs) {
if (typeof d.score !== "number" || !Number.isFinite(d.score)) {
throw rerankFailed3(
"reranker output: every document must have a finite numeric score when scores are present"
);
}
}
return docs.map(cloneDoc).sort((a, b) => b.score - a.score);
}
function validateRerankOutput(value) {
if (!Array.isArray(value)) {
throw rerankFailed3("reranker output must be an array of documents");
}
const out = [];
for (let i = 0; i < value.length; i++) {
const item = value[i];
if (item == null || typeof item !== "object") {
throw rerankFailed3(`reranker output[${i}] is not a document object`);
}
const doc = item;
if (typeof doc.id !== "string" || typeof doc.content !== "string") {
throw rerankFailed3(`reranker output[${i}] must have string id and content`);
}
if (doc.score !== void 0 && (typeof doc.score !== "number" || !Number.isFinite(doc.score))) {
throw rerankFailed3(`reranker output[${i}] has a non-finite score`);
}
out.push(cloneDoc(doc));
}
return out;
}
function createRerankedRetriever(base, options = {}) {
const candidatePool = Math.max(1, options.candidatePool ?? 20);
const topK = Math.max(1, options.topK ?? 5);
const candidatePool = resolvePositiveInt2(options.candidatePool, 20);
const topK = resolvePositiveInt2(options.topK, 5);
const rerank = options.rerank ?? bm25Rerank;
return {
async retrieve(request) {
const candidates = (await base.retrieve(request)).slice(0, candidatePool);
const reranked = await rerank({ query: request.query, documents: candidates });
return reranked.slice(0, topK);
const baseResults = await base.retrieve(request);
const candidates = baseResults.slice(0, candidatePool).map(cloneDoc);
let reranked;
try {
reranked = await rerank({ query: request.query, documents: candidates });
} catch (cause) {
if (cause instanceof RagError) throw cause;
throw rerankFailed3("reranker threw", cause);
}
const validated = validateRerankOutput(reranked);
const ordered = enforceScoreOrder2(validated);
return ordered.slice(0, topK);
}

@@ -435,8 +852,16 @@ };

}
function resolveBm25K1(value) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return 1.5;
return value;
}
function resolveBm25B(value) {
if (value === void 0 || !Number.isFinite(value) || value < 0 || value > 1) return 0.75;
return value;
}
function bm25Score(query, documents, options = {}) {
const k1 = options.k1 ?? 1.5;
const b = options.b ?? 0.75;
const k1 = resolveBm25K1(options.k1);
const b = resolveBm25B(options.b);
const qTerms = tokenize(query);
const N = documents.length;
if (N === 0 || qTerms.length === 0) return documents;
if (N === 0 || qTerms.length === 0) return documents.map(cloneDoc);
const docTerms = documents.map((d) => tokenize(d.content));

@@ -461,5 +886,8 @@ const avgdl = docTerms.reduce((acc, t) => acc + t.length, 0) / N;

const norm = 1 - b + b * (dl / (avgdl || 1));
score += idf * (f * (k1 + 1) / (f + k1 * norm));
const denom = f + k1 * norm;
if (denom === 0) continue;
const term = idf * (f * (k1 + 1) / denom);
if (Number.isFinite(term)) score += term;
}
return { ...doc, score };
return { ...cloneDoc(doc), score: Number.isFinite(score) ? score : 0 };
});

@@ -470,18 +898,34 @@ return scored.sort((a, b2) => (b2.score ?? 0) - (a.score ?? 0));

function normalize(docs) {
const scores = docs.map((d) => d.score ?? 0);
const max = Math.max(...scores, 0);
const map = /* @__PURE__ */ new Map();
const entries = [];
for (const d of docs) {
map.set(d.id, max > 0 ? (d.score ?? 0) / max : 0);
const raw = d.score;
const score = typeof raw === "number" && Number.isFinite(raw) ? raw : 0;
entries.push({ id: d.id, score });
}
const scores = entries.map((e) => e.score);
const min = scores.length > 0 ? Math.min(...scores) : 0;
const max = scores.length > 0 ? Math.max(...scores) : 0;
const range = max - min;
const map = /* @__PURE__ */ new Map();
for (const e of entries) {
if (map.has(e.id)) continue;
if (range > 0) {
map.set(e.id, (e.score - min) / range);
} else {
map.set(e.id, max === 0 ? 0 : 1);
}
}
return map;
}
function createHybridRetriever(base, options = {}) {
const vectorWeight = options.vectorWeight ?? 0.6;
const bm25Weight = options.bm25Weight ?? 0.4;
const topK = Math.max(1, options.topK ?? 5);
const candidatePool = Math.max(1, options.candidatePool ?? 20);
const { vectorWeight, bm25Weight } = resolveRelativeWeights(
options.vectorWeight,
options.bm25Weight
);
const topK = resolvePositiveInt2(options.topK, 5);
const candidatePool = resolvePositiveInt2(options.candidatePool, 20);
return {
async retrieve(request) {
const candidates = (await base.retrieve(request)).slice(0, candidatePool);
const baseResults = await base.retrieve(request);
const candidates = baseResults.slice(0, candidatePool).map(cloneDoc);
if (candidates.length === 0) return candidates;

@@ -491,6 +935,9 @@ const vectorScores = normalize(candidates);

const bm25Scores = normalize(bm25Docs);
const merged = candidates.map((d) => ({
...d,
score: vectorWeight * (vectorScores.get(d.id) ?? 0) + bm25Weight * (bm25Scores.get(d.id) ?? 0)
}));
const merged = candidates.map((d) => {
const raw = vectorWeight * (vectorScores.get(d.id) ?? 0) + bm25Weight * (bm25Scores.get(d.id) ?? 0);
return {
...d,
score: Number.isFinite(raw) ? raw : 0
};
});
merged.sort((a, b) => (b.score ?? 0) - (a.score ?? 0));

@@ -497,0 +944,0 @@ return merged.slice(0, topK);

@@ -6,2 +6,12 @@ import { AgentsKitError, generateId } from '@agentskit/core';

// src/chunker.ts
function resolveChunkSize(chunkSize) {
if (!Number.isFinite(chunkSize) || chunkSize <= 0) return Number.POSITIVE_INFINITY;
return Math.floor(chunkSize);
}
function resolveChunkOverlap(chunkOverlap, chunkSize) {
if (!Number.isFinite(chunkOverlap) || chunkOverlap < 0) return 0;
const overlap = Math.floor(chunkOverlap);
if (!Number.isFinite(chunkSize)) return 0;
return Math.min(overlap, Math.max(0, chunkSize - 1));
}
function chunkText(text, options) {

@@ -12,3 +22,4 @@ if (!text) return [];

}
const { chunkSize, chunkOverlap } = options;
const chunkSize = resolveChunkSize(options.chunkSize);
const chunkOverlap = resolveChunkOverlap(options.chunkOverlap, chunkSize);
if (text.length <= chunkSize) return [text];

@@ -37,2 +48,31 @@ const chunks = [];

// src/rag.ts
function resolvePositiveInt(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return Math.max(1, Math.floor(value));
}
function resolveThreshold(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return value;
}
function resolveChunkSize2(value, fallback) {
if (value === void 0 || !Number.isFinite(value) || value <= 0) return fallback;
return Math.floor(value);
}
function resolveChunkOverlap2(value, fallback, chunkSize) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return fallback;
return Math.min(Math.floor(value), Math.max(0, chunkSize - 1));
}
function enforceScoreOrder(docs) {
if (docs.length === 0) return docs;
const anyScore = docs.some((d) => d.score !== void 0);
if (!anyScore) return docs;
for (const d of docs) {
if (typeof d.score !== "number" || !Number.isFinite(d.score)) {
throw new TypeError(
"createRAG search: every result must have a finite numeric score when scores are present"
);
}
}
return [...docs].sort((a, b) => b.score - a.score);
}
function createRAG(config) {

@@ -42,13 +82,17 @@ const {

store,
chunkSize = 512,
chunkOverlap = 50,
split,
topK = 5,
threshold = 0
split
} = config;
const chunkSize = resolveChunkSize2(config.chunkSize, 512);
const chunkOverlap = resolveChunkOverlap2(config.chunkOverlap, 50, chunkSize);
const defaultTopK = resolvePositiveInt(config.topK, 5);
const defaultThreshold = resolveThreshold(config.threshold, 0);
async function ingest(documents) {
const vectorDocs = [];
for (const doc of documents) {
const content = doc.content;
if (!content) continue;
const docId = doc.id ?? generateId("doc");
const chunks = chunkText(doc.content, { chunkSize, chunkOverlap, split });
const source = doc.source;
const metadata = doc.metadata ? { ...doc.metadata } : void 0;
const chunks = chunkText(content, { chunkSize, chunkOverlap, split });
for (let i = 0; i < chunks.length; i++) {

@@ -62,4 +106,4 @@ const chunkId = `${docId}_chunk_${i}`;

metadata: {
...doc.metadata,
source: doc.source,
...metadata,
source,
documentId: docId,

@@ -76,7 +120,7 @@ chunkIndex: i

async function search(query, options) {
const topK = options?.topK !== void 0 ? resolvePositiveInt(options.topK, defaultTopK) : defaultTopK;
const threshold = options?.threshold !== void 0 ? resolveThreshold(options.threshold, defaultThreshold) : defaultThreshold;
const queryEmbedding = await embed(query);
return store.search(queryEmbedding, {
topK: options?.topK ?? topK,
threshold: options?.threshold ?? threshold
});
const results = await store.search(queryEmbedding, { topK, threshold });
return enforceScoreOrder(results);
}

@@ -103,8 +147,100 @@ async function retrieve(request) {

// src/loaders.ts
// src/loaders/shared.ts
function resolveMaxFiles(maxFiles, fallback = 100) {
if (maxFiles === void 0) return fallback;
if (!Number.isFinite(maxFiles)) return 0;
return Math.max(0, Math.floor(maxFiles));
}
function loadFailed(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_LOAD_FAILED,
message,
cause
});
}
function isAbortLike(err) {
if (err == null || typeof err !== "object") return false;
const name = err.name;
if (name === "AbortError") return true;
const cause = err.cause;
return cause !== void 0 && isAbortLike(cause);
}
function ensureNotAborted(signal, label) {
if (signal?.aborted) {
throw loadFailed(`${label}: aborted`, signal.reason);
}
}
function rethrowIfAbort(err, signal, label) {
ensureNotAborted(signal, label);
if (isAbortLike(err)) {
if (err instanceof RagError) throw err;
throw loadFailed(`${label}: aborted`, err);
}
}
function finishTreeLoad(label, attempted, loaded, docs) {
if (attempted > 0 && loaded === 0) {
throw loadFailed(`${label}: all eligible downloads failed`);
}
return docs;
}
async function doFetch(fetchImpl, url, init, label) {
try {
return await fetchImpl(url, init);
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) {
throw loadFailed(`${label}: aborted`, cause);
}
throw loadFailed(`${label}: network error for ${url}`, cause);
}
}
async function readResponseText(response, label) {
try {
return await response.text();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read response body`, cause);
}
}
async function readResponseJson(response, label) {
try {
return await response.json();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to parse response body`, cause);
}
}
async function readResponseArrayBuffer(response, label) {
try {
return await response.arrayBuffer();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read response body`, cause);
}
}
async function readS3Body(body, label) {
if (body == null || typeof body.transformToString !== "function") {
throw loadFailed(`${label}: missing or invalid object body`);
}
try {
return await body.transformToString();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read object body`, cause);
}
}
function encodePathSegments(path) {
return path.split("/").map((segment) => encodeURIComponent(segment)).join("/");
}
// src/loaders/documents.ts
async function loadUrl(url, options = {}) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await fetchImpl(url, { headers: options.headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadUrl ${response.status}: ${url}` });
const content = await response.text();
const response = await doFetch(fetchImpl, url, { headers: options.headers, signal: options.signal }, "loadUrl");
if (!response.ok) throw loadFailed(`loadUrl ${response.status}: ${url}`);
const content = await readResponseText(response, "loadUrl");
return [{ content, source: url, metadata: { url } }];

@@ -115,10 +251,11 @@ }

const ref = options.ref ?? "HEAD";
const url = `https://raw.githubusercontent.com/${owner}/${repo}/${ref}/${path}`;
const encodedPath = encodePathSegments(path);
const url = `https://raw.githubusercontent.com/${owner}/${repo}/${ref}/${encodedPath}`;
const headers = {};
if (options.token) headers.authorization = `Bearer ${options.token}`;
const response = await fetchImpl(url, { headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGitHubFile ${response.status}: ${url}` });
const response = await doFetch(fetchImpl, url, { headers, signal: options.signal }, "loadGitHubFile");
if (!response.ok) throw loadFailed(`loadGitHubFile ${response.status}: ${url}`);
return [
{
content: await response.text(),
content: await readResponseText(response, "loadGitHubFile"),
source: url,

@@ -132,29 +269,61 @@ metadata: { owner, repo, path, ref }

const ref = options.ref ?? "HEAD";
const url = `https://api.github.com/repos/${owner}/${repo}/git/trees/${ref}?recursive=1`;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const url = `https://api.github.com/repos/${owner}/${repo}/git/trees/${encodeURIComponent(ref)}?recursive=1`;
const headers = { accept: "application/vnd.github+json" };
if (options.token) headers.authorization = `Bearer ${options.token}`;
const response = await fetchImpl(url, { headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGitHubTree ${response.status}: ${url}` });
const tree = await response.json();
const files = (tree.tree ?? []).filter((t) => t.type === "blob").filter((t) => !options.filter || options.filter(t.path)).slice(0, options.maxFiles ?? 100);
const response = await doFetch(fetchImpl, url, { headers, signal: options.signal }, "loadGitHubTree");
if (!response.ok) throw loadFailed(`loadGitHubTree ${response.status}: ${url}`);
const tree = await readResponseJson(
response,
"loadGitHubTree"
);
const files = (tree.tree ?? []).filter((t) => t.type === "blob").filter((t) => !options.filter || options.filter(t.path)).slice(0, maxFiles);
const docs = [];
let attempted = 0;
let loaded = 0;
for (const file of files) {
const items = await loadGitHubFile(owner, repo, file.path, options);
docs.push(...items);
ensureNotAborted(options.signal, "loadGitHubTree");
attempted++;
try {
const items = await loadGitHubFile(owner, repo, file.path, options);
docs.push(...items);
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadGitHubTree");
}
}
return docs;
return finishTreeLoad("loadGitHubTree", attempted, loaded, docs);
}
async function loadNotionPage(pageId, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const url = `https://api.notion.com/v1/blocks/${pageId}/children?page_size=100`;
const response = await fetchImpl(url, {
headers: {
authorization: `Bearer ${options.token}`,
"notion-version": options.version ?? "2022-06-28"
const HEADING_PREFIX = { heading_1: "# ", heading_2: "## ", heading_3: "### " };
const blocks = [];
let cursor;
const seenCursors = /* @__PURE__ */ new Set();
while (true) {
ensureNotAborted(options.signal, "loadNotionPage");
let url = `https://api.notion.com/v1/blocks/${pageId}/children?page_size=100`;
if (cursor) url += `&start_cursor=${encodeURIComponent(cursor)}`;
const response = await doFetch(fetchImpl, url, {
headers: {
authorization: `Bearer ${options.token}`,
"notion-version": options.version ?? "2022-06-28"
},
signal: options.signal
}, "loadNotionPage");
if (!response.ok) throw loadFailed(`loadNotionPage ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadNotionPage");
for (const block of data.results ?? []) {
blocks.push(block);
}
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadNotionPage ${response.status}: ${url}` });
const data = await response.json();
const HEADING_PREFIX = { heading_1: "# ", heading_2: "## ", heading_3: "### " };
const text = (data.results ?? []).map((block) => {
if (!data.has_more) break;
const next = data.next_cursor;
if (typeof next !== "string" || next.length === 0 || seenCursors.has(next)) {
throw loadFailed("loadNotionPage: incomplete pagination (has_more without a new cursor)");
}
seenCursors.add(next);
cursor = next;
}
const text = blocks.map((block) => {
const part = block.paragraph?.rich_text ?? block.heading_1?.rich_text ?? block.heading_2?.rich_text ?? block.heading_3?.rich_text;

@@ -171,7 +340,11 @@ if (!part) return "";

const authHeader = options.authorization ?? (options.token ? `Basic ${options.token}` : void 0);
const response = await fetchImpl(url, {
headers: authHeader ? { authorization: authHeader } : {}
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadConfluencePage ${response.status}: ${url}` });
const data = await response.json();
const response = await doFetch(fetchImpl, url, {
headers: authHeader ? { authorization: authHeader } : {},
signal: options.signal
}, "loadConfluencePage");
if (!response.ok) throw loadFailed(`loadConfluencePage ${response.status}: ${url}`);
const data = await readResponseJson(
response,
"loadConfluencePage"
);
const content = data.body?.storage?.value ?? "";

@@ -183,9 +356,20 @@ return [{ content, source: `${options.baseUrl}/pages/${pageId}`, metadata: { pageId, title: data.title } }];

const url = `https://www.googleapis.com/drive/v3/files/${fileId}/export?mimeType=text/plain`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${options.accessToken}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGoogleDriveFile ${response.status}: ${url}` });
const content = await response.text();
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${options.accessToken}` },
signal: options.signal
}, "loadGoogleDriveFile");
if (!response.ok) throw loadFailed(`loadGoogleDriveFile ${response.status}: ${url}`);
const content = await readResponseText(response, "loadGoogleDriveFile");
return [{ content, source: `gdrive://${fileId}`, metadata: { fileId } }];
}
async function loadPdf(url, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await doFetch(fetchImpl, url, { signal: options.signal }, "loadPdf");
if (!response.ok) throw loadFailed(`loadPdf ${response.status}: ${url}`);
const buf = new Uint8Array(await readResponseArrayBuffer(response, "loadPdf"));
const { text, pages } = await options.parsePdf(buf);
return [{ content: text, source: url, metadata: { url, pages } }];
}
// src/loaders/s3.ts
async function loadS3(options) {

@@ -199,67 +383,108 @@ if (!options.commands) {

}
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const { ListObjectsV2Command, GetObjectCommand } = options.commands;
const docs = [];
let attempted = 0;
let loaded = 0;
let continuationToken;
const maxFiles = options.maxFiles ?? 100;
outer: while (true) {
const list = await options.client.send(new ListObjectsV2Command({
Bucket: options.bucket,
Prefix: options.prefix,
ContinuationToken: continuationToken
}));
ensureNotAborted(options.signal, "loadS3");
let list;
try {
list = await options.client.send(new ListObjectsV2Command({
Bucket: options.bucket,
Prefix: options.prefix,
ContinuationToken: continuationToken
}));
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed("loadS3: aborted", cause);
throw loadFailed(`loadS3: list failed for s3://${options.bucket}`, cause);
}
for (const obj of list.Contents ?? []) {
ensureNotAborted(options.signal, "loadS3");
const key = obj.Key;
if (!key) continue;
if (options.filter && !options.filter(key)) continue;
const get = await options.client.send(new GetObjectCommand({
Bucket: options.bucket,
Key: key
}));
const content = await get.Body?.transformToString() ?? "";
docs.push({
content,
source: `s3://${options.bucket}/${key}`,
metadata: { bucket: options.bucket, key }
});
attempted++;
try {
const get = await options.client.send(new GetObjectCommand({
Bucket: options.bucket,
Key: key
}));
const content = await readS3Body(get.Body, "loadS3");
docs.push({
content,
source: `s3://${options.bucket}/${key}`,
metadata: { bucket: options.bucket, key }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadS3");
}
if (docs.length >= maxFiles) break outer;
}
if (!list.IsTruncated) break;
continuationToken = list.NextContinuationToken;
const next = list.NextContinuationToken;
if (!next || next === continuationToken) {
throw loadFailed("loadS3: incomplete pagination (truncated without a new continuation token)");
}
continuationToken = next;
}
return docs;
return finishTreeLoad("loadS3", attempted, loaded, docs);
}
// src/loaders/cloud.ts
async function loadGcs(options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const docs = [];
const maxFiles = options.maxFiles ?? 100;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const getToken = async () => typeof options.accessToken === "string" ? options.accessToken : await options.accessToken();
let attempted = 0;
let loaded = 0;
let pageToken;
outer: while (true) {
ensureNotAborted(options.signal, "loadGcs");
const params = new URLSearchParams();
if (options.prefix) params.set("prefix", options.prefix);
if (pageToken) params.set("pageToken", pageToken);
const url = `https://storage.googleapis.com/storage/v1/b/${options.bucket}/o?${params.toString()}`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGcs ${response.status}: ${url}` });
const data = await response.json();
const url = `https://storage.googleapis.com/storage/v1/b/${encodeURIComponent(options.bucket)}/o?${params.toString()}`;
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadGcs");
if (!response.ok) throw loadFailed(`loadGcs ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadGcs");
for (const item of data.items ?? []) {
ensureNotAborted(options.signal, "loadGcs");
if (options.filter && !options.filter(item.name)) continue;
const objUrl = `https://storage.googleapis.com/storage/v1/b/${options.bucket}/o/${encodeURIComponent(item.name)}?alt=media`;
const objResponse = await fetchImpl(objUrl, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!objResponse.ok) continue;
docs.push({
content: await objResponse.text(),
source: `gs://${options.bucket}/${item.name}`,
metadata: { bucket: options.bucket, name: item.name }
});
const objUrl = `https://storage.googleapis.com/storage/v1/b/${encodeURIComponent(options.bucket)}/o/${encodeURIComponent(item.name)}?alt=media`;
attempted++;
try {
const objResponse = await doFetch(fetchImpl, objUrl, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadGcs");
if (!objResponse.ok) continue;
docs.push({
content: await readResponseText(objResponse, "loadGcs"),
source: `gs://${options.bucket}/${item.name}`,
metadata: { bucket: options.bucket, name: item.name }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadGcs");
}
if (docs.length >= maxFiles) break outer;
}
if (!data.nextPageToken) break;
pageToken = data.nextPageToken;
const next = data.nextPageToken;
if (!next) break;
if (next === pageToken) {
throw loadFailed("loadGcs: incomplete pagination (no new page token)");
}
pageToken = next;
}
return docs;
return finishTreeLoad("loadGcs", attempted, loaded, docs);
}

@@ -269,12 +494,22 @@ async function loadDropbox(options) {

const docs = [];
const maxFiles = options.maxFiles ?? 100;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const headers = { authorization: `Bearer ${options.accessToken}`, "content-type": "application/json" };
let attempted = 0;
let loaded = 0;
let cursor;
outer: while (true) {
ensureNotAborted(options.signal, "loadDropbox");
const url = cursor ? "https://api.dropboxapi.com/2/files/list_folder/continue" : "https://api.dropboxapi.com/2/files/list_folder";
const body = cursor ? { cursor } : { path: options.path ?? "", recursive: true };
const response = await fetchImpl(url, { method: "POST", headers, body: JSON.stringify(body) });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadDropbox ${response.status}: ${url}` });
const data = await response.json();
const response = await doFetch(fetchImpl, url, {
method: "POST",
headers,
body: JSON.stringify(body),
signal: options.signal
}, "loadDropbox");
if (!response.ok) throw loadFailed(`loadDropbox ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadDropbox");
for (const entry of data.entries ?? []) {
ensureNotAborted(options.signal, "loadDropbox");
if (entry[".tag"] !== "file") continue;

@@ -284,21 +519,32 @@ const path = entry.path_display ?? entry.path_lower;

if (options.filter && !options.filter(path)) continue;
const downloadResponse = await fetchImpl("https://content.dropboxapi.com/2/files/download", {
method: "POST",
headers: {
authorization: `Bearer ${options.accessToken}`,
"Dropbox-API-Arg": JSON.stringify({ path })
}
});
if (!downloadResponse.ok) continue;
docs.push({
content: await downloadResponse.text(),
source: `dropbox:${path}`,
metadata: { path }
});
attempted++;
try {
const downloadResponse = await doFetch(fetchImpl, "https://content.dropboxapi.com/2/files/download", {
method: "POST",
headers: {
authorization: `Bearer ${options.accessToken}`,
"Dropbox-API-Arg": JSON.stringify({ path })
},
signal: options.signal
}, "loadDropbox");
if (!downloadResponse.ok) continue;
docs.push({
content: await readResponseText(downloadResponse, "loadDropbox"),
source: `dropbox:${path}`,
metadata: { path }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadDropbox");
}
if (docs.length >= maxFiles) break outer;
}
if (!data.has_more) break;
cursor = data.cursor;
const next = data.cursor;
if (!next || next === cursor) {
throw loadFailed("loadDropbox: incomplete pagination (has_more without a new cursor)");
}
cursor = next;
}
return docs;
return finishTreeLoad("loadDropbox", attempted, loaded, docs);
}

@@ -308,45 +554,77 @@ async function loadOneDrive(options) {

const docs = [];
const maxFiles = options.maxFiles ?? 100;
const driveBase = options.driveId ? `https://graph.microsoft.com/v1.0/drives/${options.driveId}` : `https://graph.microsoft.com/v1.0/me/drive`;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const driveBase = options.driveId ? `https://graph.microsoft.com/v1.0/drives/${encodeURIComponent(options.driveId)}` : `https://graph.microsoft.com/v1.0/me/drive`;
const folder = options.folderItemId ? `items/${options.folderItemId}` : "root";
const getToken = async () => typeof options.accessToken === "string" ? options.accessToken : await options.accessToken();
const visitedFolders = /* @__PURE__ */ new Set();
let attempted = 0;
let loaded = 0;
async function walk(prefix) {
const url = `${driveBase}/${prefix}/children`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadOneDrive ${response.status}: ${url}` });
const data = await response.json();
for (const item of data.value ?? []) {
if (docs.length >= maxFiles) return;
if (visitedFolders.has(prefix)) return;
visitedFolders.add(prefix);
let url = `${driveBase}/${prefix}/children`;
const seenLinks = /* @__PURE__ */ new Set();
while (url) {
ensureNotAborted(options.signal, "loadOneDrive");
if (docs.length >= maxFiles) return;
if (item.folder) {
await walk(`items/${item.id}`);
if (seenLinks.has(url)) {
throw loadFailed("loadOneDrive: incomplete pagination (repeated @odata.nextLink)");
}
seenLinks.add(url);
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadOneDrive");
if (!response.ok) throw loadFailed(`loadOneDrive ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadOneDrive");
for (const item of data.value ?? []) {
ensureNotAborted(options.signal, "loadOneDrive");
if (docs.length >= maxFiles) return;
if (item.folder) {
await walk(`items/${item.id}`);
continue;
}
if (!item.file) continue;
if (options.filter && !options.filter(item.name)) continue;
const downloadUrl = item["@microsoft.graph.downloadUrl"];
if (!downloadUrl) continue;
attempted++;
try {
const fileResponse = await doFetch(fetchImpl, downloadUrl, { signal: options.signal }, "loadOneDrive");
if (!fileResponse.ok) continue;
docs.push({
content: await readResponseText(fileResponse, "loadOneDrive"),
source: `onedrive:${item.id}`,
metadata: { id: item.id, name: item.name, mimeType: item.file.mimeType }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadOneDrive");
}
}
const next = data["@odata.nextLink"];
if (!next) {
url = void 0;
continue;
}
if (!item.file) continue;
if (options.filter && !options.filter(item.name)) continue;
const downloadUrl = item["@microsoft.graph.downloadUrl"];
if (!downloadUrl) continue;
const fileResponse = await fetchImpl(downloadUrl);
if (!fileResponse.ok) continue;
docs.push({
content: await fileResponse.text(),
source: `onedrive:${item.id}`,
metadata: { id: item.id, name: item.name, mimeType: item.file.mimeType }
});
if (seenLinks.has(next) || next === url) {
throw loadFailed("loadOneDrive: incomplete pagination (repeated @odata.nextLink)");
}
url = next;
}
}
await walk(folder);
return docs;
return finishTreeLoad("loadOneDrive", attempted, loaded, docs);
}
async function loadPdf(url, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await fetchImpl(url);
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadPdf ${response.status}: ${url}` });
const buf = new Uint8Array(await response.arrayBuffer());
const { text, pages } = await options.parsePdf(buf);
return [{ content: text, source: url, metadata: { url, pages } }];
}
// src/rerankers/voyage.ts
function rerankFailed(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function voyageReranker(options) {

@@ -357,23 +635,52 @@ const fetchImpl = options.fetch ?? globalThis.fetch;

if (documents.length === 0) return documents;
const response = await fetchImpl("https://api.voyageai.com/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
})
});
let response;
try {
response = await fetchImpl("https://api.voyageai.com/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
}),
signal: options.signal
});
} catch (cause) {
throw rerankFailed("voyage rerank: network error", cause);
}
if (!response.ok) {
const text = await response.text();
throw new RagError({ code: RagErrorCodes.AK_RAG_RERANK_FAILED, message: `voyage rerank: ${response.status} ${text.slice(0, 200)}` });
let text = "";
try {
text = await response.text();
} catch (cause) {
throw rerankFailed(`voyage rerank: ${response.status} (failed to read error body)`, cause);
}
throw rerankFailed(`voyage rerank: ${response.status} ${text.slice(0, 200)}`);
}
const data = await response.json();
const ranked = (data.data ?? []).sort((a, b) => b.relevance_score - a.relevance_score).map((r) => {
let data;
try {
data = await response.json();
} catch (cause) {
throw rerankFailed("voyage rerank: invalid JSON response", cause);
}
if (!Array.isArray(data.data)) {
throw rerankFailed("voyage rerank: data must be an array");
}
const ranked = [];
for (let i = 0; i < data.data.length; i++) {
const r = data.data[i];
if (r == null || typeof r !== "object" || !Number.isInteger(r.index) || r.index < 0 || r.index >= documents.length || typeof r.relevance_score !== "number" || !Number.isFinite(r.relevance_score)) {
throw rerankFailed(`voyage rerank: malformed result at index ${i}`);
}
const doc = documents[r.index];
return { ...doc, score: r.relevance_score };
});
ranked.push({
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0,
score: r.relevance_score
});
}
ranked.sort((a, b) => b.score - a.score);
return ranked;

@@ -384,2 +691,9 @@ };

// src/rerankers/jina.ts
function rerankFailed2(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function jinaReranker(options) {

@@ -390,23 +704,52 @@ const fetchImpl = options.fetch ?? globalThis.fetch;

if (documents.length === 0) return documents;
const response = await fetchImpl("https://api.jina.ai/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
})
});
let response;
try {
response = await fetchImpl("https://api.jina.ai/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
}),
signal: options.signal
});
} catch (cause) {
throw rerankFailed2("jina rerank: network error", cause);
}
if (!response.ok) {
const text = await response.text();
throw new RagError({ code: RagErrorCodes.AK_RAG_RERANK_FAILED, message: `jina rerank: ${response.status} ${text.slice(0, 200)}` });
let text = "";
try {
text = await response.text();
} catch (cause) {
throw rerankFailed2(`jina rerank: ${response.status} (failed to read error body)`, cause);
}
throw rerankFailed2(`jina rerank: ${response.status} ${text.slice(0, 200)}`);
}
const data = await response.json();
const ranked = (data.results ?? []).sort((a, b) => b.relevance_score - a.relevance_score).map((r) => {
let data;
try {
data = await response.json();
} catch (cause) {
throw rerankFailed2("jina rerank: invalid JSON response", cause);
}
if (!Array.isArray(data.results)) {
throw rerankFailed2("jina rerank: results must be an array");
}
const ranked = [];
for (let i = 0; i < data.results.length; i++) {
const r = data.results[i];
if (r == null || typeof r !== "object" || !Number.isInteger(r.index) || r.index < 0 || r.index >= documents.length || typeof r.relevance_score !== "number" || !Number.isFinite(r.relevance_score)) {
throw rerankFailed2(`jina rerank: malformed result at index ${i}`);
}
const doc = documents[r.index];
return { ...doc, score: r.relevance_score };
});
ranked.push({
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0,
score: r.relevance_score
});
}
ranked.sort((a, b) => b.score - a.score);
return ranked;

@@ -417,11 +760,85 @@ };

// src/rerank.ts
function resolvePositiveInt2(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return Math.max(1, Math.floor(value));
}
function resolveWeight(value, fallback) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return fallback;
return value;
}
function resolveRelativeWeights(vectorWeight, bm25Weight) {
const v = resolveWeight(vectorWeight, 0.6);
const b = resolveWeight(bm25Weight, 0.4);
if (v === 0 && b === 0) return { vectorWeight: 0.5, bm25Weight: 0.5 };
const m = Math.max(v, b);
const vN = v / m;
const bN = b / m;
const sum = vN + bN;
return { vectorWeight: vN / sum, bm25Weight: bN / sum };
}
function cloneDoc(doc) {
return {
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0
};
}
function rerankFailed3(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function enforceScoreOrder2(docs) {
if (docs.length === 0) return docs;
const anyScore = docs.some((d) => d.score !== void 0);
if (!anyScore) return docs.map(cloneDoc);
for (const d of docs) {
if (typeof d.score !== "number" || !Number.isFinite(d.score)) {
throw rerankFailed3(
"reranker output: every document must have a finite numeric score when scores are present"
);
}
}
return docs.map(cloneDoc).sort((a, b) => b.score - a.score);
}
function validateRerankOutput(value) {
if (!Array.isArray(value)) {
throw rerankFailed3("reranker output must be an array of documents");
}
const out = [];
for (let i = 0; i < value.length; i++) {
const item = value[i];
if (item == null || typeof item !== "object") {
throw rerankFailed3(`reranker output[${i}] is not a document object`);
}
const doc = item;
if (typeof doc.id !== "string" || typeof doc.content !== "string") {
throw rerankFailed3(`reranker output[${i}] must have string id and content`);
}
if (doc.score !== void 0 && (typeof doc.score !== "number" || !Number.isFinite(doc.score))) {
throw rerankFailed3(`reranker output[${i}] has a non-finite score`);
}
out.push(cloneDoc(doc));
}
return out;
}
function createRerankedRetriever(base, options = {}) {
const candidatePool = Math.max(1, options.candidatePool ?? 20);
const topK = Math.max(1, options.topK ?? 5);
const candidatePool = resolvePositiveInt2(options.candidatePool, 20);
const topK = resolvePositiveInt2(options.topK, 5);
const rerank = options.rerank ?? bm25Rerank;
return {
async retrieve(request) {
const candidates = (await base.retrieve(request)).slice(0, candidatePool);
const reranked = await rerank({ query: request.query, documents: candidates });
return reranked.slice(0, topK);
const baseResults = await base.retrieve(request);
const candidates = baseResults.slice(0, candidatePool).map(cloneDoc);
let reranked;
try {
reranked = await rerank({ query: request.query, documents: candidates });
} catch (cause) {
if (cause instanceof RagError) throw cause;
throw rerankFailed3("reranker threw", cause);
}
const validated = validateRerankOutput(reranked);
const ordered = enforceScoreOrder2(validated);
return ordered.slice(0, topK);
}

@@ -433,8 +850,16 @@ };

}
function resolveBm25K1(value) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return 1.5;
return value;
}
function resolveBm25B(value) {
if (value === void 0 || !Number.isFinite(value) || value < 0 || value > 1) return 0.75;
return value;
}
function bm25Score(query, documents, options = {}) {
const k1 = options.k1 ?? 1.5;
const b = options.b ?? 0.75;
const k1 = resolveBm25K1(options.k1);
const b = resolveBm25B(options.b);
const qTerms = tokenize(query);
const N = documents.length;
if (N === 0 || qTerms.length === 0) return documents;
if (N === 0 || qTerms.length === 0) return documents.map(cloneDoc);
const docTerms = documents.map((d) => tokenize(d.content));

@@ -459,5 +884,8 @@ const avgdl = docTerms.reduce((acc, t) => acc + t.length, 0) / N;

const norm = 1 - b + b * (dl / (avgdl || 1));
score += idf * (f * (k1 + 1) / (f + k1 * norm));
const denom = f + k1 * norm;
if (denom === 0) continue;
const term = idf * (f * (k1 + 1) / denom);
if (Number.isFinite(term)) score += term;
}
return { ...doc, score };
return { ...cloneDoc(doc), score: Number.isFinite(score) ? score : 0 };
});

@@ -468,18 +896,34 @@ return scored.sort((a, b2) => (b2.score ?? 0) - (a.score ?? 0));

function normalize(docs) {
const scores = docs.map((d) => d.score ?? 0);
const max = Math.max(...scores, 0);
const map = /* @__PURE__ */ new Map();
const entries = [];
for (const d of docs) {
map.set(d.id, max > 0 ? (d.score ?? 0) / max : 0);
const raw = d.score;
const score = typeof raw === "number" && Number.isFinite(raw) ? raw : 0;
entries.push({ id: d.id, score });
}
const scores = entries.map((e) => e.score);
const min = scores.length > 0 ? Math.min(...scores) : 0;
const max = scores.length > 0 ? Math.max(...scores) : 0;
const range = max - min;
const map = /* @__PURE__ */ new Map();
for (const e of entries) {
if (map.has(e.id)) continue;
if (range > 0) {
map.set(e.id, (e.score - min) / range);
} else {
map.set(e.id, max === 0 ? 0 : 1);
}
}
return map;
}
function createHybridRetriever(base, options = {}) {
const vectorWeight = options.vectorWeight ?? 0.6;
const bm25Weight = options.bm25Weight ?? 0.4;
const topK = Math.max(1, options.topK ?? 5);
const candidatePool = Math.max(1, options.candidatePool ?? 20);
const { vectorWeight, bm25Weight } = resolveRelativeWeights(
options.vectorWeight,
options.bm25Weight
);
const topK = resolvePositiveInt2(options.topK, 5);
const candidatePool = resolvePositiveInt2(options.candidatePool, 20);
return {
async retrieve(request) {
const candidates = (await base.retrieve(request)).slice(0, candidatePool);
const baseResults = await base.retrieve(request);
const candidates = baseResults.slice(0, candidatePool).map(cloneDoc);
if (candidates.length === 0) return candidates;

@@ -489,6 +933,9 @@ const vectorScores = normalize(candidates);

const bm25Scores = normalize(bm25Docs);
const merged = candidates.map((d) => ({
...d,
score: vectorWeight * (vectorScores.get(d.id) ?? 0) + bm25Weight * (bm25Scores.get(d.id) ?? 0)
}));
const merged = candidates.map((d) => {
const raw = vectorWeight * (vectorScores.get(d.id) ?? 0) + bm25Weight * (bm25Scores.get(d.id) ?? 0);
return {
...d,
score: Number.isFinite(raw) ? raw : 0
};
});
merged.sort((a, b) => (b.score ?? 0) - (a.score ?? 0));

@@ -495,0 +942,0 @@ return merged.slice(0, topK);

@@ -8,2 +8,12 @@ 'use strict';

// src/chunker.ts
function resolveChunkSize(chunkSize) {
if (!Number.isFinite(chunkSize) || chunkSize <= 0) return Number.POSITIVE_INFINITY;
return Math.floor(chunkSize);
}
function resolveChunkOverlap(chunkOverlap, chunkSize) {
if (!Number.isFinite(chunkOverlap) || chunkOverlap < 0) return 0;
const overlap = Math.floor(chunkOverlap);
if (!Number.isFinite(chunkSize)) return 0;
return Math.min(overlap, Math.max(0, chunkSize - 1));
}
function chunkText(text, options) {

@@ -14,3 +24,4 @@ if (!text) return [];

}
const { chunkSize, chunkOverlap } = options;
const chunkSize = resolveChunkSize(options.chunkSize);
const chunkOverlap = resolveChunkOverlap(options.chunkOverlap, chunkSize);
if (text.length <= chunkSize) return [text];

@@ -39,2 +50,31 @@ const chunks = [];

// src/rag.ts
function resolvePositiveInt(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return Math.max(1, Math.floor(value));
}
function resolveThreshold(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return value;
}
function resolveChunkSize2(value, fallback) {
if (value === void 0 || !Number.isFinite(value) || value <= 0) return fallback;
return Math.floor(value);
}
function resolveChunkOverlap2(value, fallback, chunkSize) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return fallback;
return Math.min(Math.floor(value), Math.max(0, chunkSize - 1));
}
function enforceScoreOrder(docs) {
if (docs.length === 0) return docs;
const anyScore = docs.some((d) => d.score !== void 0);
if (!anyScore) return docs;
for (const d of docs) {
if (typeof d.score !== "number" || !Number.isFinite(d.score)) {
throw new TypeError(
"createRAG search: every result must have a finite numeric score when scores are present"
);
}
}
return [...docs].sort((a, b) => b.score - a.score);
}
function createRAG(config) {

@@ -44,13 +84,17 @@ const {

store,
chunkSize = 512,
chunkOverlap = 50,
split,
topK = 5,
threshold = 0
split
} = config;
const chunkSize = resolveChunkSize2(config.chunkSize, 512);
const chunkOverlap = resolveChunkOverlap2(config.chunkOverlap, 50, chunkSize);
const defaultTopK = resolvePositiveInt(config.topK, 5);
const defaultThreshold = resolveThreshold(config.threshold, 0);
async function ingest(documents) {
const vectorDocs = [];
for (const doc of documents) {
const content = doc.content;
if (!content) continue;
const docId = doc.id ?? core.generateId("doc");
const chunks = chunkText(doc.content, { chunkSize, chunkOverlap, split });
const source = doc.source;
const metadata = doc.metadata ? { ...doc.metadata } : void 0;
const chunks = chunkText(content, { chunkSize, chunkOverlap, split });
for (let i = 0; i < chunks.length; i++) {

@@ -64,4 +108,4 @@ const chunkId = `${docId}_chunk_${i}`;

metadata: {
...doc.metadata,
source: doc.source,
...metadata,
source,
documentId: docId,

@@ -78,7 +122,7 @@ chunkIndex: i

async function search(query, options) {
const topK = options?.topK !== void 0 ? resolvePositiveInt(options.topK, defaultTopK) : defaultTopK;
const threshold = options?.threshold !== void 0 ? resolveThreshold(options.threshold, defaultThreshold) : defaultThreshold;
const queryEmbedding = await embed(query);
return store.search(queryEmbedding, {
topK: options?.topK ?? topK,
threshold: options?.threshold ?? threshold
});
const results = await store.search(queryEmbedding, { topK, threshold });
return enforceScoreOrder(results);
}

@@ -105,8 +149,100 @@ async function retrieve(request) {

// src/loaders.ts
// src/loaders/shared.ts
function resolveMaxFiles(maxFiles, fallback = 100) {
if (maxFiles === void 0) return fallback;
if (!Number.isFinite(maxFiles)) return 0;
return Math.max(0, Math.floor(maxFiles));
}
function loadFailed(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_LOAD_FAILED,
message,
cause
});
}
function isAbortLike(err) {
if (err == null || typeof err !== "object") return false;
const name = err.name;
if (name === "AbortError") return true;
const cause = err.cause;
return cause !== void 0 && isAbortLike(cause);
}
function ensureNotAborted(signal, label) {
if (signal?.aborted) {
throw loadFailed(`${label}: aborted`, signal.reason);
}
}
function rethrowIfAbort(err, signal, label) {
ensureNotAborted(signal, label);
if (isAbortLike(err)) {
if (err instanceof RagError) throw err;
throw loadFailed(`${label}: aborted`, err);
}
}
function finishTreeLoad(label, attempted, loaded, docs) {
if (attempted > 0 && loaded === 0) {
throw loadFailed(`${label}: all eligible downloads failed`);
}
return docs;
}
async function doFetch(fetchImpl, url, init, label) {
try {
return await fetchImpl(url, init);
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) {
throw loadFailed(`${label}: aborted`, cause);
}
throw loadFailed(`${label}: network error for ${url}`, cause);
}
}
async function readResponseText(response, label) {
try {
return await response.text();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read response body`, cause);
}
}
async function readResponseJson(response, label) {
try {
return await response.json();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to parse response body`, cause);
}
}
async function readResponseArrayBuffer(response, label) {
try {
return await response.arrayBuffer();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read response body`, cause);
}
}
async function readS3Body(body, label) {
if (body == null || typeof body.transformToString !== "function") {
throw loadFailed(`${label}: missing or invalid object body`);
}
try {
return await body.transformToString();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read object body`, cause);
}
}
function encodePathSegments(path) {
return path.split("/").map((segment) => encodeURIComponent(segment)).join("/");
}
// src/loaders/documents.ts
async function loadUrl(url, options = {}) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await fetchImpl(url, { headers: options.headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadUrl ${response.status}: ${url}` });
const content = await response.text();
const response = await doFetch(fetchImpl, url, { headers: options.headers, signal: options.signal }, "loadUrl");
if (!response.ok) throw loadFailed(`loadUrl ${response.status}: ${url}`);
const content = await readResponseText(response, "loadUrl");
return [{ content, source: url, metadata: { url } }];

@@ -117,10 +253,11 @@ }

const ref = options.ref ?? "HEAD";
const url = `https://raw.githubusercontent.com/${owner}/${repo}/${ref}/${path}`;
const encodedPath = encodePathSegments(path);
const url = `https://raw.githubusercontent.com/${owner}/${repo}/${ref}/${encodedPath}`;
const headers = {};
if (options.token) headers.authorization = `Bearer ${options.token}`;
const response = await fetchImpl(url, { headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGitHubFile ${response.status}: ${url}` });
const response = await doFetch(fetchImpl, url, { headers, signal: options.signal }, "loadGitHubFile");
if (!response.ok) throw loadFailed(`loadGitHubFile ${response.status}: ${url}`);
return [
{
content: await response.text(),
content: await readResponseText(response, "loadGitHubFile"),
source: url,

@@ -134,29 +271,61 @@ metadata: { owner, repo, path, ref }

const ref = options.ref ?? "HEAD";
const url = `https://api.github.com/repos/${owner}/${repo}/git/trees/${ref}?recursive=1`;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const url = `https://api.github.com/repos/${owner}/${repo}/git/trees/${encodeURIComponent(ref)}?recursive=1`;
const headers = { accept: "application/vnd.github+json" };
if (options.token) headers.authorization = `Bearer ${options.token}`;
const response = await fetchImpl(url, { headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGitHubTree ${response.status}: ${url}` });
const tree = await response.json();
const files = (tree.tree ?? []).filter((t) => t.type === "blob").filter((t) => !options.filter || options.filter(t.path)).slice(0, options.maxFiles ?? 100);
const response = await doFetch(fetchImpl, url, { headers, signal: options.signal }, "loadGitHubTree");
if (!response.ok) throw loadFailed(`loadGitHubTree ${response.status}: ${url}`);
const tree = await readResponseJson(
response,
"loadGitHubTree"
);
const files = (tree.tree ?? []).filter((t) => t.type === "blob").filter((t) => !options.filter || options.filter(t.path)).slice(0, maxFiles);
const docs = [];
let attempted = 0;
let loaded = 0;
for (const file of files) {
const items = await loadGitHubFile(owner, repo, file.path, options);
docs.push(...items);
ensureNotAborted(options.signal, "loadGitHubTree");
attempted++;
try {
const items = await loadGitHubFile(owner, repo, file.path, options);
docs.push(...items);
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadGitHubTree");
}
}
return docs;
return finishTreeLoad("loadGitHubTree", attempted, loaded, docs);
}
async function loadNotionPage(pageId, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const url = `https://api.notion.com/v1/blocks/${pageId}/children?page_size=100`;
const response = await fetchImpl(url, {
headers: {
authorization: `Bearer ${options.token}`,
"notion-version": options.version ?? "2022-06-28"
const HEADING_PREFIX = { heading_1: "# ", heading_2: "## ", heading_3: "### " };
const blocks = [];
let cursor;
const seenCursors = /* @__PURE__ */ new Set();
while (true) {
ensureNotAborted(options.signal, "loadNotionPage");
let url = `https://api.notion.com/v1/blocks/${pageId}/children?page_size=100`;
if (cursor) url += `&start_cursor=${encodeURIComponent(cursor)}`;
const response = await doFetch(fetchImpl, url, {
headers: {
authorization: `Bearer ${options.token}`,
"notion-version": options.version ?? "2022-06-28"
},
signal: options.signal
}, "loadNotionPage");
if (!response.ok) throw loadFailed(`loadNotionPage ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadNotionPage");
for (const block of data.results ?? []) {
blocks.push(block);
}
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadNotionPage ${response.status}: ${url}` });
const data = await response.json();
const HEADING_PREFIX = { heading_1: "# ", heading_2: "## ", heading_3: "### " };
const text = (data.results ?? []).map((block) => {
if (!data.has_more) break;
const next = data.next_cursor;
if (typeof next !== "string" || next.length === 0 || seenCursors.has(next)) {
throw loadFailed("loadNotionPage: incomplete pagination (has_more without a new cursor)");
}
seenCursors.add(next);
cursor = next;
}
const text = blocks.map((block) => {
const part = block.paragraph?.rich_text ?? block.heading_1?.rich_text ?? block.heading_2?.rich_text ?? block.heading_3?.rich_text;

@@ -173,7 +342,11 @@ if (!part) return "";

const authHeader = options.authorization ?? (options.token ? `Basic ${options.token}` : void 0);
const response = await fetchImpl(url, {
headers: authHeader ? { authorization: authHeader } : {}
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadConfluencePage ${response.status}: ${url}` });
const data = await response.json();
const response = await doFetch(fetchImpl, url, {
headers: authHeader ? { authorization: authHeader } : {},
signal: options.signal
}, "loadConfluencePage");
if (!response.ok) throw loadFailed(`loadConfluencePage ${response.status}: ${url}`);
const data = await readResponseJson(
response,
"loadConfluencePage"
);
const content = data.body?.storage?.value ?? "";

@@ -185,9 +358,20 @@ return [{ content, source: `${options.baseUrl}/pages/${pageId}`, metadata: { pageId, title: data.title } }];

const url = `https://www.googleapis.com/drive/v3/files/${fileId}/export?mimeType=text/plain`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${options.accessToken}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGoogleDriveFile ${response.status}: ${url}` });
const content = await response.text();
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${options.accessToken}` },
signal: options.signal
}, "loadGoogleDriveFile");
if (!response.ok) throw loadFailed(`loadGoogleDriveFile ${response.status}: ${url}`);
const content = await readResponseText(response, "loadGoogleDriveFile");
return [{ content, source: `gdrive://${fileId}`, metadata: { fileId } }];
}
async function loadPdf(url, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await doFetch(fetchImpl, url, { signal: options.signal }, "loadPdf");
if (!response.ok) throw loadFailed(`loadPdf ${response.status}: ${url}`);
const buf = new Uint8Array(await readResponseArrayBuffer(response, "loadPdf"));
const { text, pages } = await options.parsePdf(buf);
return [{ content: text, source: url, metadata: { url, pages } }];
}
// src/loaders/s3.ts
async function loadS3(options) {

@@ -201,67 +385,108 @@ if (!options.commands) {

}
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const { ListObjectsV2Command, GetObjectCommand } = options.commands;
const docs = [];
let attempted = 0;
let loaded = 0;
let continuationToken;
const maxFiles = options.maxFiles ?? 100;
outer: while (true) {
const list = await options.client.send(new ListObjectsV2Command({
Bucket: options.bucket,
Prefix: options.prefix,
ContinuationToken: continuationToken
}));
ensureNotAborted(options.signal, "loadS3");
let list;
try {
list = await options.client.send(new ListObjectsV2Command({
Bucket: options.bucket,
Prefix: options.prefix,
ContinuationToken: continuationToken
}));
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed("loadS3: aborted", cause);
throw loadFailed(`loadS3: list failed for s3://${options.bucket}`, cause);
}
for (const obj of list.Contents ?? []) {
ensureNotAborted(options.signal, "loadS3");
const key = obj.Key;
if (!key) continue;
if (options.filter && !options.filter(key)) continue;
const get = await options.client.send(new GetObjectCommand({
Bucket: options.bucket,
Key: key
}));
const content = await get.Body?.transformToString() ?? "";
docs.push({
content,
source: `s3://${options.bucket}/${key}`,
metadata: { bucket: options.bucket, key }
});
attempted++;
try {
const get = await options.client.send(new GetObjectCommand({
Bucket: options.bucket,
Key: key
}));
const content = await readS3Body(get.Body, "loadS3");
docs.push({
content,
source: `s3://${options.bucket}/${key}`,
metadata: { bucket: options.bucket, key }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadS3");
}
if (docs.length >= maxFiles) break outer;
}
if (!list.IsTruncated) break;
continuationToken = list.NextContinuationToken;
const next = list.NextContinuationToken;
if (!next || next === continuationToken) {
throw loadFailed("loadS3: incomplete pagination (truncated without a new continuation token)");
}
continuationToken = next;
}
return docs;
return finishTreeLoad("loadS3", attempted, loaded, docs);
}
// src/loaders/cloud.ts
async function loadGcs(options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const docs = [];
const maxFiles = options.maxFiles ?? 100;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const getToken = async () => typeof options.accessToken === "string" ? options.accessToken : await options.accessToken();
let attempted = 0;
let loaded = 0;
let pageToken;
outer: while (true) {
ensureNotAborted(options.signal, "loadGcs");
const params = new URLSearchParams();
if (options.prefix) params.set("prefix", options.prefix);
if (pageToken) params.set("pageToken", pageToken);
const url = `https://storage.googleapis.com/storage/v1/b/${options.bucket}/o?${params.toString()}`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGcs ${response.status}: ${url}` });
const data = await response.json();
const url = `https://storage.googleapis.com/storage/v1/b/${encodeURIComponent(options.bucket)}/o?${params.toString()}`;
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadGcs");
if (!response.ok) throw loadFailed(`loadGcs ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadGcs");
for (const item of data.items ?? []) {
ensureNotAborted(options.signal, "loadGcs");
if (options.filter && !options.filter(item.name)) continue;
const objUrl = `https://storage.googleapis.com/storage/v1/b/${options.bucket}/o/${encodeURIComponent(item.name)}?alt=media`;
const objResponse = await fetchImpl(objUrl, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!objResponse.ok) continue;
docs.push({
content: await objResponse.text(),
source: `gs://${options.bucket}/${item.name}`,
metadata: { bucket: options.bucket, name: item.name }
});
const objUrl = `https://storage.googleapis.com/storage/v1/b/${encodeURIComponent(options.bucket)}/o/${encodeURIComponent(item.name)}?alt=media`;
attempted++;
try {
const objResponse = await doFetch(fetchImpl, objUrl, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadGcs");
if (!objResponse.ok) continue;
docs.push({
content: await readResponseText(objResponse, "loadGcs"),
source: `gs://${options.bucket}/${item.name}`,
metadata: { bucket: options.bucket, name: item.name }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadGcs");
}
if (docs.length >= maxFiles) break outer;
}
if (!data.nextPageToken) break;
pageToken = data.nextPageToken;
const next = data.nextPageToken;
if (!next) break;
if (next === pageToken) {
throw loadFailed("loadGcs: incomplete pagination (no new page token)");
}
pageToken = next;
}
return docs;
return finishTreeLoad("loadGcs", attempted, loaded, docs);
}

@@ -271,12 +496,22 @@ async function loadDropbox(options) {

const docs = [];
const maxFiles = options.maxFiles ?? 100;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const headers = { authorization: `Bearer ${options.accessToken}`, "content-type": "application/json" };
let attempted = 0;
let loaded = 0;
let cursor;
outer: while (true) {
ensureNotAborted(options.signal, "loadDropbox");
const url = cursor ? "https://api.dropboxapi.com/2/files/list_folder/continue" : "https://api.dropboxapi.com/2/files/list_folder";
const body = cursor ? { cursor } : { path: options.path ?? "", recursive: true };
const response = await fetchImpl(url, { method: "POST", headers, body: JSON.stringify(body) });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadDropbox ${response.status}: ${url}` });
const data = await response.json();
const response = await doFetch(fetchImpl, url, {
method: "POST",
headers,
body: JSON.stringify(body),
signal: options.signal
}, "loadDropbox");
if (!response.ok) throw loadFailed(`loadDropbox ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadDropbox");
for (const entry of data.entries ?? []) {
ensureNotAborted(options.signal, "loadDropbox");
if (entry[".tag"] !== "file") continue;

@@ -286,21 +521,32 @@ const path = entry.path_display ?? entry.path_lower;

if (options.filter && !options.filter(path)) continue;
const downloadResponse = await fetchImpl("https://content.dropboxapi.com/2/files/download", {
method: "POST",
headers: {
authorization: `Bearer ${options.accessToken}`,
"Dropbox-API-Arg": JSON.stringify({ path })
}
});
if (!downloadResponse.ok) continue;
docs.push({
content: await downloadResponse.text(),
source: `dropbox:${path}`,
metadata: { path }
});
attempted++;
try {
const downloadResponse = await doFetch(fetchImpl, "https://content.dropboxapi.com/2/files/download", {
method: "POST",
headers: {
authorization: `Bearer ${options.accessToken}`,
"Dropbox-API-Arg": JSON.stringify({ path })
},
signal: options.signal
}, "loadDropbox");
if (!downloadResponse.ok) continue;
docs.push({
content: await readResponseText(downloadResponse, "loadDropbox"),
source: `dropbox:${path}`,
metadata: { path }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadDropbox");
}
if (docs.length >= maxFiles) break outer;
}
if (!data.has_more) break;
cursor = data.cursor;
const next = data.cursor;
if (!next || next === cursor) {
throw loadFailed("loadDropbox: incomplete pagination (has_more without a new cursor)");
}
cursor = next;
}
return docs;
return finishTreeLoad("loadDropbox", attempted, loaded, docs);
}

@@ -310,45 +556,77 @@ async function loadOneDrive(options) {

const docs = [];
const maxFiles = options.maxFiles ?? 100;
const driveBase = options.driveId ? `https://graph.microsoft.com/v1.0/drives/${options.driveId}` : `https://graph.microsoft.com/v1.0/me/drive`;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const driveBase = options.driveId ? `https://graph.microsoft.com/v1.0/drives/${encodeURIComponent(options.driveId)}` : `https://graph.microsoft.com/v1.0/me/drive`;
const folder = options.folderItemId ? `items/${options.folderItemId}` : "root";
const getToken = async () => typeof options.accessToken === "string" ? options.accessToken : await options.accessToken();
const visitedFolders = /* @__PURE__ */ new Set();
let attempted = 0;
let loaded = 0;
async function walk(prefix) {
const url = `${driveBase}/${prefix}/children`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadOneDrive ${response.status}: ${url}` });
const data = await response.json();
for (const item of data.value ?? []) {
if (docs.length >= maxFiles) return;
if (visitedFolders.has(prefix)) return;
visitedFolders.add(prefix);
let url = `${driveBase}/${prefix}/children`;
const seenLinks = /* @__PURE__ */ new Set();
while (url) {
ensureNotAborted(options.signal, "loadOneDrive");
if (docs.length >= maxFiles) return;
if (item.folder) {
await walk(`items/${item.id}`);
if (seenLinks.has(url)) {
throw loadFailed("loadOneDrive: incomplete pagination (repeated @odata.nextLink)");
}
seenLinks.add(url);
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadOneDrive");
if (!response.ok) throw loadFailed(`loadOneDrive ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadOneDrive");
for (const item of data.value ?? []) {
ensureNotAborted(options.signal, "loadOneDrive");
if (docs.length >= maxFiles) return;
if (item.folder) {
await walk(`items/${item.id}`);
continue;
}
if (!item.file) continue;
if (options.filter && !options.filter(item.name)) continue;
const downloadUrl = item["@microsoft.graph.downloadUrl"];
if (!downloadUrl) continue;
attempted++;
try {
const fileResponse = await doFetch(fetchImpl, downloadUrl, { signal: options.signal }, "loadOneDrive");
if (!fileResponse.ok) continue;
docs.push({
content: await readResponseText(fileResponse, "loadOneDrive"),
source: `onedrive:${item.id}`,
metadata: { id: item.id, name: item.name, mimeType: item.file.mimeType }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadOneDrive");
}
}
const next = data["@odata.nextLink"];
if (!next) {
url = void 0;
continue;
}
if (!item.file) continue;
if (options.filter && !options.filter(item.name)) continue;
const downloadUrl = item["@microsoft.graph.downloadUrl"];
if (!downloadUrl) continue;
const fileResponse = await fetchImpl(downloadUrl);
if (!fileResponse.ok) continue;
docs.push({
content: await fileResponse.text(),
source: `onedrive:${item.id}`,
metadata: { id: item.id, name: item.name, mimeType: item.file.mimeType }
});
if (seenLinks.has(next) || next === url) {
throw loadFailed("loadOneDrive: incomplete pagination (repeated @odata.nextLink)");
}
url = next;
}
}
await walk(folder);
return docs;
return finishTreeLoad("loadOneDrive", attempted, loaded, docs);
}
async function loadPdf(url, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await fetchImpl(url);
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadPdf ${response.status}: ${url}` });
const buf = new Uint8Array(await response.arrayBuffer());
const { text, pages } = await options.parsePdf(buf);
return [{ content: text, source: url, metadata: { url, pages } }];
}
// src/rerankers/voyage.ts
function rerankFailed(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function voyageReranker(options) {

@@ -359,23 +637,52 @@ const fetchImpl = options.fetch ?? globalThis.fetch;

if (documents.length === 0) return documents;
const response = await fetchImpl("https://api.voyageai.com/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
})
});
let response;
try {
response = await fetchImpl("https://api.voyageai.com/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
}),
signal: options.signal
});
} catch (cause) {
throw rerankFailed("voyage rerank: network error", cause);
}
if (!response.ok) {
const text = await response.text();
throw new RagError({ code: RagErrorCodes.AK_RAG_RERANK_FAILED, message: `voyage rerank: ${response.status} ${text.slice(0, 200)}` });
let text = "";
try {
text = await response.text();
} catch (cause) {
throw rerankFailed(`voyage rerank: ${response.status} (failed to read error body)`, cause);
}
throw rerankFailed(`voyage rerank: ${response.status} ${text.slice(0, 200)}`);
}
const data = await response.json();
const ranked = (data.data ?? []).sort((a, b) => b.relevance_score - a.relevance_score).map((r) => {
let data;
try {
data = await response.json();
} catch (cause) {
throw rerankFailed("voyage rerank: invalid JSON response", cause);
}
if (!Array.isArray(data.data)) {
throw rerankFailed("voyage rerank: data must be an array");
}
const ranked = [];
for (let i = 0; i < data.data.length; i++) {
const r = data.data[i];
if (r == null || typeof r !== "object" || !Number.isInteger(r.index) || r.index < 0 || r.index >= documents.length || typeof r.relevance_score !== "number" || !Number.isFinite(r.relevance_score)) {
throw rerankFailed(`voyage rerank: malformed result at index ${i}`);
}
const doc = documents[r.index];
return { ...doc, score: r.relevance_score };
});
ranked.push({
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0,
score: r.relevance_score
});
}
ranked.sort((a, b) => b.score - a.score);
return ranked;

@@ -386,2 +693,9 @@ };

// src/rerankers/jina.ts
function rerankFailed2(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function jinaReranker(options) {

@@ -392,23 +706,52 @@ const fetchImpl = options.fetch ?? globalThis.fetch;

if (documents.length === 0) return documents;
const response = await fetchImpl("https://api.jina.ai/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
})
});
let response;
try {
response = await fetchImpl("https://api.jina.ai/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
}),
signal: options.signal
});
} catch (cause) {
throw rerankFailed2("jina rerank: network error", cause);
}
if (!response.ok) {
const text = await response.text();
throw new RagError({ code: RagErrorCodes.AK_RAG_RERANK_FAILED, message: `jina rerank: ${response.status} ${text.slice(0, 200)}` });
let text = "";
try {
text = await response.text();
} catch (cause) {
throw rerankFailed2(`jina rerank: ${response.status} (failed to read error body)`, cause);
}
throw rerankFailed2(`jina rerank: ${response.status} ${text.slice(0, 200)}`);
}
const data = await response.json();
const ranked = (data.results ?? []).sort((a, b) => b.relevance_score - a.relevance_score).map((r) => {
let data;
try {
data = await response.json();
} catch (cause) {
throw rerankFailed2("jina rerank: invalid JSON response", cause);
}
if (!Array.isArray(data.results)) {
throw rerankFailed2("jina rerank: results must be an array");
}
const ranked = [];
for (let i = 0; i < data.results.length; i++) {
const r = data.results[i];
if (r == null || typeof r !== "object" || !Number.isInteger(r.index) || r.index < 0 || r.index >= documents.length || typeof r.relevance_score !== "number" || !Number.isFinite(r.relevance_score)) {
throw rerankFailed2(`jina rerank: malformed result at index ${i}`);
}
const doc = documents[r.index];
return { ...doc, score: r.relevance_score };
});
ranked.push({
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0,
score: r.relevance_score
});
}
ranked.sort((a, b) => b.score - a.score);
return ranked;

@@ -419,11 +762,85 @@ };

// src/rerank.ts
function resolvePositiveInt2(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return Math.max(1, Math.floor(value));
}
function resolveWeight(value, fallback) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return fallback;
return value;
}
function resolveRelativeWeights(vectorWeight, bm25Weight) {
const v = resolveWeight(vectorWeight, 0.6);
const b = resolveWeight(bm25Weight, 0.4);
if (v === 0 && b === 0) return { vectorWeight: 0.5, bm25Weight: 0.5 };
const m = Math.max(v, b);
const vN = v / m;
const bN = b / m;
const sum = vN + bN;
return { vectorWeight: vN / sum, bm25Weight: bN / sum };
}
function cloneDoc(doc) {
return {
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0
};
}
function rerankFailed3(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function enforceScoreOrder2(docs) {
if (docs.length === 0) return docs;
const anyScore = docs.some((d) => d.score !== void 0);
if (!anyScore) return docs.map(cloneDoc);
for (const d of docs) {
if (typeof d.score !== "number" || !Number.isFinite(d.score)) {
throw rerankFailed3(
"reranker output: every document must have a finite numeric score when scores are present"
);
}
}
return docs.map(cloneDoc).sort((a, b) => b.score - a.score);
}
function validateRerankOutput(value) {
if (!Array.isArray(value)) {
throw rerankFailed3("reranker output must be an array of documents");
}
const out = [];
for (let i = 0; i < value.length; i++) {
const item = value[i];
if (item == null || typeof item !== "object") {
throw rerankFailed3(`reranker output[${i}] is not a document object`);
}
const doc = item;
if (typeof doc.id !== "string" || typeof doc.content !== "string") {
throw rerankFailed3(`reranker output[${i}] must have string id and content`);
}
if (doc.score !== void 0 && (typeof doc.score !== "number" || !Number.isFinite(doc.score))) {
throw rerankFailed3(`reranker output[${i}] has a non-finite score`);
}
out.push(cloneDoc(doc));
}
return out;
}
function createRerankedRetriever(base, options = {}) {
const candidatePool = Math.max(1, options.candidatePool ?? 20);
const topK = Math.max(1, options.topK ?? 5);
const candidatePool = resolvePositiveInt2(options.candidatePool, 20);
const topK = resolvePositiveInt2(options.topK, 5);
const rerank = options.rerank ?? bm25Rerank;
return {
async retrieve(request) {
const candidates = (await base.retrieve(request)).slice(0, candidatePool);
const reranked = await rerank({ query: request.query, documents: candidates });
return reranked.slice(0, topK);
const baseResults = await base.retrieve(request);
const candidates = baseResults.slice(0, candidatePool).map(cloneDoc);
let reranked;
try {
reranked = await rerank({ query: request.query, documents: candidates });
} catch (cause) {
if (cause instanceof RagError) throw cause;
throw rerankFailed3("reranker threw", cause);
}
const validated = validateRerankOutput(reranked);
const ordered = enforceScoreOrder2(validated);
return ordered.slice(0, topK);
}

@@ -435,8 +852,16 @@ };

}
function resolveBm25K1(value) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return 1.5;
return value;
}
function resolveBm25B(value) {
if (value === void 0 || !Number.isFinite(value) || value < 0 || value > 1) return 0.75;
return value;
}
function bm25Score(query, documents, options = {}) {
const k1 = options.k1 ?? 1.5;
const b = options.b ?? 0.75;
const k1 = resolveBm25K1(options.k1);
const b = resolveBm25B(options.b);
const qTerms = tokenize(query);
const N = documents.length;
if (N === 0 || qTerms.length === 0) return documents;
if (N === 0 || qTerms.length === 0) return documents.map(cloneDoc);
const docTerms = documents.map((d) => tokenize(d.content));

@@ -461,5 +886,8 @@ const avgdl = docTerms.reduce((acc, t) => acc + t.length, 0) / N;

const norm = 1 - b + b * (dl / (avgdl || 1));
score += idf * (f * (k1 + 1) / (f + k1 * norm));
const denom = f + k1 * norm;
if (denom === 0) continue;
const term = idf * (f * (k1 + 1) / denom);
if (Number.isFinite(term)) score += term;
}
return { ...doc, score };
return { ...cloneDoc(doc), score: Number.isFinite(score) ? score : 0 };
});

@@ -470,18 +898,34 @@ return scored.sort((a, b2) => (b2.score ?? 0) - (a.score ?? 0));

function normalize(docs) {
const scores = docs.map((d) => d.score ?? 0);
const max = Math.max(...scores, 0);
const map = /* @__PURE__ */ new Map();
const entries = [];
for (const d of docs) {
map.set(d.id, max > 0 ? (d.score ?? 0) / max : 0);
const raw = d.score;
const score = typeof raw === "number" && Number.isFinite(raw) ? raw : 0;
entries.push({ id: d.id, score });
}
const scores = entries.map((e) => e.score);
const min = scores.length > 0 ? Math.min(...scores) : 0;
const max = scores.length > 0 ? Math.max(...scores) : 0;
const range = max - min;
const map = /* @__PURE__ */ new Map();
for (const e of entries) {
if (map.has(e.id)) continue;
if (range > 0) {
map.set(e.id, (e.score - min) / range);
} else {
map.set(e.id, max === 0 ? 0 : 1);
}
}
return map;
}
function createHybridRetriever(base, options = {}) {
const vectorWeight = options.vectorWeight ?? 0.6;
const bm25Weight = options.bm25Weight ?? 0.4;
const topK = Math.max(1, options.topK ?? 5);
const candidatePool = Math.max(1, options.candidatePool ?? 20);
const { vectorWeight, bm25Weight } = resolveRelativeWeights(
options.vectorWeight,
options.bm25Weight
);
const topK = resolvePositiveInt2(options.topK, 5);
const candidatePool = resolvePositiveInt2(options.candidatePool, 20);
return {
async retrieve(request) {
const candidates = (await base.retrieve(request)).slice(0, candidatePool);
const baseResults = await base.retrieve(request);
const candidates = baseResults.slice(0, candidatePool).map(cloneDoc);
if (candidates.length === 0) return candidates;

@@ -491,6 +935,9 @@ const vectorScores = normalize(candidates);

const bm25Scores = normalize(bm25Docs);
const merged = candidates.map((d) => ({
...d,
score: vectorWeight * (vectorScores.get(d.id) ?? 0) + bm25Weight * (bm25Scores.get(d.id) ?? 0)
}));
const merged = candidates.map((d) => {
const raw = vectorWeight * (vectorScores.get(d.id) ?? 0) + bm25Weight * (bm25Scores.get(d.id) ?? 0);
return {
...d,
score: Number.isFinite(raw) ? raw : 0
};
});
merged.sort((a, b) => (b.score ?? 0) - (a.score ?? 0));

@@ -497,0 +944,0 @@ return merged.slice(0, topK);

@@ -58,9 +58,8 @@ import { Retriever, RetrieverRequest, RetrievedDocument, EmbedFn, VectorMemory, AgentsKitError } from '@agentskit/core';

/**
* Document loaders: small async functions returning `InputDocument[]` ready to
* pipe into `RAG.ingest`. Framework-agnostic; each accepts a custom `fetch`.
*/
interface LoaderOptions {
fetch?: typeof globalThis.fetch;
/** Optional abort signal forwarded to underlying HTTP calls when supported. */
signal?: AbortSignal;
}
interface UrlLoaderOptions extends LoaderOptions {

@@ -99,2 +98,17 @@ headers?: Record<string, string>;

declare function loadGoogleDriveFile(fileId: string, options: DriveLoaderOptions): Promise<InputDocument[]>;
interface PdfLoaderOptions extends LoaderOptions {
parsePdf: (bytes: Uint8Array) => Promise<{
text: string;
pages?: number;
}> | {
text: string;
pages?: number;
};
}
/**
* PDF loader — parser is BYO so native deps stay out of the bundle.
* Fetch bytes at `url`, hand to `parsePdf`, wrap in `InputDocument`.
*/
declare function loadPdf(url: string, options: PdfLoaderOptions): Promise<InputDocument[]>;
interface S3LikeClient {

@@ -133,2 +147,3 @@ send(command: {

}
interface GcsLoaderOptions extends LoaderOptions {

@@ -163,16 +178,2 @@ bucket: string;

declare function loadOneDrive(options: OneDriveLoaderOptions): Promise<InputDocument[]>;
interface PdfLoaderOptions extends LoaderOptions {
parsePdf: (bytes: Uint8Array) => Promise<{
text: string;
pages?: number;
}> | {
text: string;
pages?: number;
};
}
/**
* PDF loader — parser is BYO so native deps stay out of the bundle.
* Fetch bytes at `url`, hand to `parsePdf`, wrap in `InputDocument`.
*/
declare function loadPdf(url: string, options: PdfLoaderOptions): Promise<InputDocument[]>;

@@ -200,5 +201,5 @@ type RerankFn = (input: {

interface BM25Options {
/** Free term-saturation knob. Default 1.5. */
/** Term-frequency saturation (k1 ≥ 0). Default 1.5; invalid → default. */
k1?: number;
/** Length-normalization weight. Default 0.75. */
/** Length-normalization weight in [0, 1]. Default 0.75; invalid → default. */
b?: number;

@@ -208,4 +209,4 @@ }

* Score a set of documents against a query using classic BM25.
* Returns the input documents augmented with a `.score` field,
* sorted by score descending.
* Returns new document objects with a finite `.score` field, sorted descending.
* Input documents are never mutated.
*/

@@ -228,3 +229,4 @@ declare function bm25Score(query: string, documents: RetrievedDocument[], options?: BM25Options): RetrievedDocument[];

* over the same candidate pool. Final score is a weighted sum of the
* two normalized scores.
* two min-max-normalized scores using a finite relative weight pair
* that sums to 1 (both zero → 0.5/0.5).
*/

@@ -239,2 +241,4 @@ declare function createHybridRetriever(base: Retriever, options?: HybridRetrieverOptions): Retriever;

fetch?: typeof globalThis.fetch;
/** Optional abort signal forwarded to the underlying HTTP request. */
signal?: AbortSignal;
}

@@ -252,2 +256,4 @@ /**

fetch?: typeof globalThis.fetch;
/** Optional abort signal forwarded to the underlying HTTP request. */
signal?: AbortSignal;
}

@@ -254,0 +260,0 @@ /**

@@ -58,9 +58,8 @@ import { Retriever, RetrieverRequest, RetrievedDocument, EmbedFn, VectorMemory, AgentsKitError } from '@agentskit/core';

/**
* Document loaders: small async functions returning `InputDocument[]` ready to
* pipe into `RAG.ingest`. Framework-agnostic; each accepts a custom `fetch`.
*/
interface LoaderOptions {
fetch?: typeof globalThis.fetch;
/** Optional abort signal forwarded to underlying HTTP calls when supported. */
signal?: AbortSignal;
}
interface UrlLoaderOptions extends LoaderOptions {

@@ -99,2 +98,17 @@ headers?: Record<string, string>;

declare function loadGoogleDriveFile(fileId: string, options: DriveLoaderOptions): Promise<InputDocument[]>;
interface PdfLoaderOptions extends LoaderOptions {
parsePdf: (bytes: Uint8Array) => Promise<{
text: string;
pages?: number;
}> | {
text: string;
pages?: number;
};
}
/**
* PDF loader — parser is BYO so native deps stay out of the bundle.
* Fetch bytes at `url`, hand to `parsePdf`, wrap in `InputDocument`.
*/
declare function loadPdf(url: string, options: PdfLoaderOptions): Promise<InputDocument[]>;
interface S3LikeClient {

@@ -133,2 +147,3 @@ send(command: {

}
interface GcsLoaderOptions extends LoaderOptions {

@@ -163,16 +178,2 @@ bucket: string;

declare function loadOneDrive(options: OneDriveLoaderOptions): Promise<InputDocument[]>;
interface PdfLoaderOptions extends LoaderOptions {
parsePdf: (bytes: Uint8Array) => Promise<{
text: string;
pages?: number;
}> | {
text: string;
pages?: number;
};
}
/**
* PDF loader — parser is BYO so native deps stay out of the bundle.
* Fetch bytes at `url`, hand to `parsePdf`, wrap in `InputDocument`.
*/
declare function loadPdf(url: string, options: PdfLoaderOptions): Promise<InputDocument[]>;

@@ -200,5 +201,5 @@ type RerankFn = (input: {

interface BM25Options {
/** Free term-saturation knob. Default 1.5. */
/** Term-frequency saturation (k1 ≥ 0). Default 1.5; invalid → default. */
k1?: number;
/** Length-normalization weight. Default 0.75. */
/** Length-normalization weight in [0, 1]. Default 0.75; invalid → default. */
b?: number;

@@ -208,4 +209,4 @@ }

* Score a set of documents against a query using classic BM25.
* Returns the input documents augmented with a `.score` field,
* sorted by score descending.
* Returns new document objects with a finite `.score` field, sorted descending.
* Input documents are never mutated.
*/

@@ -228,3 +229,4 @@ declare function bm25Score(query: string, documents: RetrievedDocument[], options?: BM25Options): RetrievedDocument[];

* over the same candidate pool. Final score is a weighted sum of the
* two normalized scores.
* two min-max-normalized scores using a finite relative weight pair
* that sums to 1 (both zero → 0.5/0.5).
*/

@@ -239,2 +241,4 @@ declare function createHybridRetriever(base: Retriever, options?: HybridRetrieverOptions): Retriever;

fetch?: typeof globalThis.fetch;
/** Optional abort signal forwarded to the underlying HTTP request. */
signal?: AbortSignal;
}

@@ -252,2 +256,4 @@ /**

fetch?: typeof globalThis.fetch;
/** Optional abort signal forwarded to the underlying HTTP request. */
signal?: AbortSignal;
}

@@ -254,0 +260,0 @@ /**

@@ -1,5 +0,34 @@

import { chunkText } from './chunk-UNFVK5RA.js';
export { chunkText } from './chunk-UNFVK5RA.js';
import { chunkText } from './chunk-K6CR4AHC.js';
export { chunkText } from './chunk-K6CR4AHC.js';
import { AgentsKitError, generateId } from '@agentskit/core';
function resolvePositiveInt(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return Math.max(1, Math.floor(value));
}
function resolveThreshold(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return value;
}
function resolveChunkSize(value, fallback) {
if (value === void 0 || !Number.isFinite(value) || value <= 0) return fallback;
return Math.floor(value);
}
function resolveChunkOverlap(value, fallback, chunkSize) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return fallback;
return Math.min(Math.floor(value), Math.max(0, chunkSize - 1));
}
function enforceScoreOrder(docs) {
if (docs.length === 0) return docs;
const anyScore = docs.some((d) => d.score !== void 0);
if (!anyScore) return docs;
for (const d of docs) {
if (typeof d.score !== "number" || !Number.isFinite(d.score)) {
throw new TypeError(
"createRAG search: every result must have a finite numeric score when scores are present"
);
}
}
return [...docs].sort((a, b) => b.score - a.score);
}
function createRAG(config) {

@@ -9,13 +38,17 @@ const {

store,
chunkSize = 512,
chunkOverlap = 50,
split,
topK = 5,
threshold = 0
split
} = config;
const chunkSize = resolveChunkSize(config.chunkSize, 512);
const chunkOverlap = resolveChunkOverlap(config.chunkOverlap, 50, chunkSize);
const defaultTopK = resolvePositiveInt(config.topK, 5);
const defaultThreshold = resolveThreshold(config.threshold, 0);
async function ingest(documents) {
const vectorDocs = [];
for (const doc of documents) {
const content = doc.content;
if (!content) continue;
const docId = doc.id ?? generateId("doc");
const chunks = chunkText(doc.content, { chunkSize, chunkOverlap, split });
const source = doc.source;
const metadata = doc.metadata ? { ...doc.metadata } : void 0;
const chunks = chunkText(content, { chunkSize, chunkOverlap, split });
for (let i = 0; i < chunks.length; i++) {

@@ -29,4 +62,4 @@ const chunkId = `${docId}_chunk_${i}`;

metadata: {
...doc.metadata,
source: doc.source,
...metadata,
source,
documentId: docId,

@@ -43,7 +76,7 @@ chunkIndex: i

async function search(query, options) {
const topK = options?.topK !== void 0 ? resolvePositiveInt(options.topK, defaultTopK) : defaultTopK;
const threshold = options?.threshold !== void 0 ? resolveThreshold(options.threshold, defaultThreshold) : defaultThreshold;
const queryEmbedding = await embed(query);
return store.search(queryEmbedding, {
topK: options?.topK ?? topK,
threshold: options?.threshold ?? threshold
});
const results = await store.search(queryEmbedding, { topK, threshold });
return enforceScoreOrder(results);
}

@@ -70,8 +103,100 @@ async function retrieve(request) {

// src/loaders.ts
// src/loaders/shared.ts
function resolveMaxFiles(maxFiles, fallback = 100) {
if (maxFiles === void 0) return fallback;
if (!Number.isFinite(maxFiles)) return 0;
return Math.max(0, Math.floor(maxFiles));
}
function loadFailed(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_LOAD_FAILED,
message,
cause
});
}
function isAbortLike(err) {
if (err == null || typeof err !== "object") return false;
const name = err.name;
if (name === "AbortError") return true;
const cause = err.cause;
return cause !== void 0 && isAbortLike(cause);
}
function ensureNotAborted(signal, label) {
if (signal?.aborted) {
throw loadFailed(`${label}: aborted`, signal.reason);
}
}
function rethrowIfAbort(err, signal, label) {
ensureNotAborted(signal, label);
if (isAbortLike(err)) {
if (err instanceof RagError) throw err;
throw loadFailed(`${label}: aborted`, err);
}
}
function finishTreeLoad(label, attempted, loaded, docs) {
if (attempted > 0 && loaded === 0) {
throw loadFailed(`${label}: all eligible downloads failed`);
}
return docs;
}
async function doFetch(fetchImpl, url, init, label) {
try {
return await fetchImpl(url, init);
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) {
throw loadFailed(`${label}: aborted`, cause);
}
throw loadFailed(`${label}: network error for ${url}`, cause);
}
}
async function readResponseText(response, label) {
try {
return await response.text();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read response body`, cause);
}
}
async function readResponseJson(response, label) {
try {
return await response.json();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to parse response body`, cause);
}
}
async function readResponseArrayBuffer(response, label) {
try {
return await response.arrayBuffer();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read response body`, cause);
}
}
async function readS3Body(body, label) {
if (body == null || typeof body.transformToString !== "function") {
throw loadFailed(`${label}: missing or invalid object body`);
}
try {
return await body.transformToString();
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed(`${label}: aborted`, cause);
throw loadFailed(`${label}: failed to read object body`, cause);
}
}
function encodePathSegments(path) {
return path.split("/").map((segment) => encodeURIComponent(segment)).join("/");
}
// src/loaders/documents.ts
async function loadUrl(url, options = {}) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await fetchImpl(url, { headers: options.headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadUrl ${response.status}: ${url}` });
const content = await response.text();
const response = await doFetch(fetchImpl, url, { headers: options.headers, signal: options.signal }, "loadUrl");
if (!response.ok) throw loadFailed(`loadUrl ${response.status}: ${url}`);
const content = await readResponseText(response, "loadUrl");
return [{ content, source: url, metadata: { url } }];

@@ -82,10 +207,11 @@ }

const ref = options.ref ?? "HEAD";
const url = `https://raw.githubusercontent.com/${owner}/${repo}/${ref}/${path}`;
const encodedPath = encodePathSegments(path);
const url = `https://raw.githubusercontent.com/${owner}/${repo}/${ref}/${encodedPath}`;
const headers = {};
if (options.token) headers.authorization = `Bearer ${options.token}`;
const response = await fetchImpl(url, { headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGitHubFile ${response.status}: ${url}` });
const response = await doFetch(fetchImpl, url, { headers, signal: options.signal }, "loadGitHubFile");
if (!response.ok) throw loadFailed(`loadGitHubFile ${response.status}: ${url}`);
return [
{
content: await response.text(),
content: await readResponseText(response, "loadGitHubFile"),
source: url,

@@ -99,29 +225,61 @@ metadata: { owner, repo, path, ref }

const ref = options.ref ?? "HEAD";
const url = `https://api.github.com/repos/${owner}/${repo}/git/trees/${ref}?recursive=1`;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const url = `https://api.github.com/repos/${owner}/${repo}/git/trees/${encodeURIComponent(ref)}?recursive=1`;
const headers = { accept: "application/vnd.github+json" };
if (options.token) headers.authorization = `Bearer ${options.token}`;
const response = await fetchImpl(url, { headers });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGitHubTree ${response.status}: ${url}` });
const tree = await response.json();
const files = (tree.tree ?? []).filter((t) => t.type === "blob").filter((t) => !options.filter || options.filter(t.path)).slice(0, options.maxFiles ?? 100);
const response = await doFetch(fetchImpl, url, { headers, signal: options.signal }, "loadGitHubTree");
if (!response.ok) throw loadFailed(`loadGitHubTree ${response.status}: ${url}`);
const tree = await readResponseJson(
response,
"loadGitHubTree"
);
const files = (tree.tree ?? []).filter((t) => t.type === "blob").filter((t) => !options.filter || options.filter(t.path)).slice(0, maxFiles);
const docs = [];
let attempted = 0;
let loaded = 0;
for (const file of files) {
const items = await loadGitHubFile(owner, repo, file.path, options);
docs.push(...items);
ensureNotAborted(options.signal, "loadGitHubTree");
attempted++;
try {
const items = await loadGitHubFile(owner, repo, file.path, options);
docs.push(...items);
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadGitHubTree");
}
}
return docs;
return finishTreeLoad("loadGitHubTree", attempted, loaded, docs);
}
async function loadNotionPage(pageId, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const url = `https://api.notion.com/v1/blocks/${pageId}/children?page_size=100`;
const response = await fetchImpl(url, {
headers: {
authorization: `Bearer ${options.token}`,
"notion-version": options.version ?? "2022-06-28"
const HEADING_PREFIX = { heading_1: "# ", heading_2: "## ", heading_3: "### " };
const blocks = [];
let cursor;
const seenCursors = /* @__PURE__ */ new Set();
while (true) {
ensureNotAborted(options.signal, "loadNotionPage");
let url = `https://api.notion.com/v1/blocks/${pageId}/children?page_size=100`;
if (cursor) url += `&start_cursor=${encodeURIComponent(cursor)}`;
const response = await doFetch(fetchImpl, url, {
headers: {
authorization: `Bearer ${options.token}`,
"notion-version": options.version ?? "2022-06-28"
},
signal: options.signal
}, "loadNotionPage");
if (!response.ok) throw loadFailed(`loadNotionPage ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadNotionPage");
for (const block of data.results ?? []) {
blocks.push(block);
}
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadNotionPage ${response.status}: ${url}` });
const data = await response.json();
const HEADING_PREFIX = { heading_1: "# ", heading_2: "## ", heading_3: "### " };
const text = (data.results ?? []).map((block) => {
if (!data.has_more) break;
const next = data.next_cursor;
if (typeof next !== "string" || next.length === 0 || seenCursors.has(next)) {
throw loadFailed("loadNotionPage: incomplete pagination (has_more without a new cursor)");
}
seenCursors.add(next);
cursor = next;
}
const text = blocks.map((block) => {
const part = block.paragraph?.rich_text ?? block.heading_1?.rich_text ?? block.heading_2?.rich_text ?? block.heading_3?.rich_text;

@@ -138,7 +296,11 @@ if (!part) return "";

const authHeader = options.authorization ?? (options.token ? `Basic ${options.token}` : void 0);
const response = await fetchImpl(url, {
headers: authHeader ? { authorization: authHeader } : {}
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadConfluencePage ${response.status}: ${url}` });
const data = await response.json();
const response = await doFetch(fetchImpl, url, {
headers: authHeader ? { authorization: authHeader } : {},
signal: options.signal
}, "loadConfluencePage");
if (!response.ok) throw loadFailed(`loadConfluencePage ${response.status}: ${url}`);
const data = await readResponseJson(
response,
"loadConfluencePage"
);
const content = data.body?.storage?.value ?? "";

@@ -150,9 +312,20 @@ return [{ content, source: `${options.baseUrl}/pages/${pageId}`, metadata: { pageId, title: data.title } }];

const url = `https://www.googleapis.com/drive/v3/files/${fileId}/export?mimeType=text/plain`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${options.accessToken}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGoogleDriveFile ${response.status}: ${url}` });
const content = await response.text();
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${options.accessToken}` },
signal: options.signal
}, "loadGoogleDriveFile");
if (!response.ok) throw loadFailed(`loadGoogleDriveFile ${response.status}: ${url}`);
const content = await readResponseText(response, "loadGoogleDriveFile");
return [{ content, source: `gdrive://${fileId}`, metadata: { fileId } }];
}
async function loadPdf(url, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await doFetch(fetchImpl, url, { signal: options.signal }, "loadPdf");
if (!response.ok) throw loadFailed(`loadPdf ${response.status}: ${url}`);
const buf = new Uint8Array(await readResponseArrayBuffer(response, "loadPdf"));
const { text, pages } = await options.parsePdf(buf);
return [{ content: text, source: url, metadata: { url, pages } }];
}
// src/loaders/s3.ts
async function loadS3(options) {

@@ -166,67 +339,108 @@ if (!options.commands) {

}
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const { ListObjectsV2Command, GetObjectCommand } = options.commands;
const docs = [];
let attempted = 0;
let loaded = 0;
let continuationToken;
const maxFiles = options.maxFiles ?? 100;
outer: while (true) {
const list = await options.client.send(new ListObjectsV2Command({
Bucket: options.bucket,
Prefix: options.prefix,
ContinuationToken: continuationToken
}));
ensureNotAborted(options.signal, "loadS3");
let list;
try {
list = await options.client.send(new ListObjectsV2Command({
Bucket: options.bucket,
Prefix: options.prefix,
ContinuationToken: continuationToken
}));
} catch (cause) {
if (cause instanceof RagError) throw cause;
if (isAbortLike(cause)) throw loadFailed("loadS3: aborted", cause);
throw loadFailed(`loadS3: list failed for s3://${options.bucket}`, cause);
}
for (const obj of list.Contents ?? []) {
ensureNotAborted(options.signal, "loadS3");
const key = obj.Key;
if (!key) continue;
if (options.filter && !options.filter(key)) continue;
const get = await options.client.send(new GetObjectCommand({
Bucket: options.bucket,
Key: key
}));
const content = await get.Body?.transformToString() ?? "";
docs.push({
content,
source: `s3://${options.bucket}/${key}`,
metadata: { bucket: options.bucket, key }
});
attempted++;
try {
const get = await options.client.send(new GetObjectCommand({
Bucket: options.bucket,
Key: key
}));
const content = await readS3Body(get.Body, "loadS3");
docs.push({
content,
source: `s3://${options.bucket}/${key}`,
metadata: { bucket: options.bucket, key }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadS3");
}
if (docs.length >= maxFiles) break outer;
}
if (!list.IsTruncated) break;
continuationToken = list.NextContinuationToken;
const next = list.NextContinuationToken;
if (!next || next === continuationToken) {
throw loadFailed("loadS3: incomplete pagination (truncated without a new continuation token)");
}
continuationToken = next;
}
return docs;
return finishTreeLoad("loadS3", attempted, loaded, docs);
}
// src/loaders/cloud.ts
async function loadGcs(options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const docs = [];
const maxFiles = options.maxFiles ?? 100;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const getToken = async () => typeof options.accessToken === "string" ? options.accessToken : await options.accessToken();
let attempted = 0;
let loaded = 0;
let pageToken;
outer: while (true) {
ensureNotAborted(options.signal, "loadGcs");
const params = new URLSearchParams();
if (options.prefix) params.set("prefix", options.prefix);
if (pageToken) params.set("pageToken", pageToken);
const url = `https://storage.googleapis.com/storage/v1/b/${options.bucket}/o?${params.toString()}`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadGcs ${response.status}: ${url}` });
const data = await response.json();
const url = `https://storage.googleapis.com/storage/v1/b/${encodeURIComponent(options.bucket)}/o?${params.toString()}`;
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadGcs");
if (!response.ok) throw loadFailed(`loadGcs ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadGcs");
for (const item of data.items ?? []) {
ensureNotAborted(options.signal, "loadGcs");
if (options.filter && !options.filter(item.name)) continue;
const objUrl = `https://storage.googleapis.com/storage/v1/b/${options.bucket}/o/${encodeURIComponent(item.name)}?alt=media`;
const objResponse = await fetchImpl(objUrl, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!objResponse.ok) continue;
docs.push({
content: await objResponse.text(),
source: `gs://${options.bucket}/${item.name}`,
metadata: { bucket: options.bucket, name: item.name }
});
const objUrl = `https://storage.googleapis.com/storage/v1/b/${encodeURIComponent(options.bucket)}/o/${encodeURIComponent(item.name)}?alt=media`;
attempted++;
try {
const objResponse = await doFetch(fetchImpl, objUrl, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadGcs");
if (!objResponse.ok) continue;
docs.push({
content: await readResponseText(objResponse, "loadGcs"),
source: `gs://${options.bucket}/${item.name}`,
metadata: { bucket: options.bucket, name: item.name }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadGcs");
}
if (docs.length >= maxFiles) break outer;
}
if (!data.nextPageToken) break;
pageToken = data.nextPageToken;
const next = data.nextPageToken;
if (!next) break;
if (next === pageToken) {
throw loadFailed("loadGcs: incomplete pagination (no new page token)");
}
pageToken = next;
}
return docs;
return finishTreeLoad("loadGcs", attempted, loaded, docs);
}

@@ -236,12 +450,22 @@ async function loadDropbox(options) {

const docs = [];
const maxFiles = options.maxFiles ?? 100;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const headers = { authorization: `Bearer ${options.accessToken}`, "content-type": "application/json" };
let attempted = 0;
let loaded = 0;
let cursor;
outer: while (true) {
ensureNotAborted(options.signal, "loadDropbox");
const url = cursor ? "https://api.dropboxapi.com/2/files/list_folder/continue" : "https://api.dropboxapi.com/2/files/list_folder";
const body = cursor ? { cursor } : { path: options.path ?? "", recursive: true };
const response = await fetchImpl(url, { method: "POST", headers, body: JSON.stringify(body) });
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadDropbox ${response.status}: ${url}` });
const data = await response.json();
const response = await doFetch(fetchImpl, url, {
method: "POST",
headers,
body: JSON.stringify(body),
signal: options.signal
}, "loadDropbox");
if (!response.ok) throw loadFailed(`loadDropbox ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadDropbox");
for (const entry of data.entries ?? []) {
ensureNotAborted(options.signal, "loadDropbox");
if (entry[".tag"] !== "file") continue;

@@ -251,21 +475,32 @@ const path = entry.path_display ?? entry.path_lower;

if (options.filter && !options.filter(path)) continue;
const downloadResponse = await fetchImpl("https://content.dropboxapi.com/2/files/download", {
method: "POST",
headers: {
authorization: `Bearer ${options.accessToken}`,
"Dropbox-API-Arg": JSON.stringify({ path })
}
});
if (!downloadResponse.ok) continue;
docs.push({
content: await downloadResponse.text(),
source: `dropbox:${path}`,
metadata: { path }
});
attempted++;
try {
const downloadResponse = await doFetch(fetchImpl, "https://content.dropboxapi.com/2/files/download", {
method: "POST",
headers: {
authorization: `Bearer ${options.accessToken}`,
"Dropbox-API-Arg": JSON.stringify({ path })
},
signal: options.signal
}, "loadDropbox");
if (!downloadResponse.ok) continue;
docs.push({
content: await readResponseText(downloadResponse, "loadDropbox"),
source: `dropbox:${path}`,
metadata: { path }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadDropbox");
}
if (docs.length >= maxFiles) break outer;
}
if (!data.has_more) break;
cursor = data.cursor;
const next = data.cursor;
if (!next || next === cursor) {
throw loadFailed("loadDropbox: incomplete pagination (has_more without a new cursor)");
}
cursor = next;
}
return docs;
return finishTreeLoad("loadDropbox", attempted, loaded, docs);
}

@@ -275,45 +510,77 @@ async function loadOneDrive(options) {

const docs = [];
const maxFiles = options.maxFiles ?? 100;
const driveBase = options.driveId ? `https://graph.microsoft.com/v1.0/drives/${options.driveId}` : `https://graph.microsoft.com/v1.0/me/drive`;
const maxFiles = resolveMaxFiles(options.maxFiles);
if (maxFiles === 0) return [];
const driveBase = options.driveId ? `https://graph.microsoft.com/v1.0/drives/${encodeURIComponent(options.driveId)}` : `https://graph.microsoft.com/v1.0/me/drive`;
const folder = options.folderItemId ? `items/${options.folderItemId}` : "root";
const getToken = async () => typeof options.accessToken === "string" ? options.accessToken : await options.accessToken();
const visitedFolders = /* @__PURE__ */ new Set();
let attempted = 0;
let loaded = 0;
async function walk(prefix) {
const url = `${driveBase}/${prefix}/children`;
const response = await fetchImpl(url, {
headers: { authorization: `Bearer ${await getToken()}` }
});
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadOneDrive ${response.status}: ${url}` });
const data = await response.json();
for (const item of data.value ?? []) {
if (docs.length >= maxFiles) return;
if (visitedFolders.has(prefix)) return;
visitedFolders.add(prefix);
let url = `${driveBase}/${prefix}/children`;
const seenLinks = /* @__PURE__ */ new Set();
while (url) {
ensureNotAborted(options.signal, "loadOneDrive");
if (docs.length >= maxFiles) return;
if (item.folder) {
await walk(`items/${item.id}`);
if (seenLinks.has(url)) {
throw loadFailed("loadOneDrive: incomplete pagination (repeated @odata.nextLink)");
}
seenLinks.add(url);
const response = await doFetch(fetchImpl, url, {
headers: { authorization: `Bearer ${await getToken()}` },
signal: options.signal
}, "loadOneDrive");
if (!response.ok) throw loadFailed(`loadOneDrive ${response.status}: ${url}`);
const data = await readResponseJson(response, "loadOneDrive");
for (const item of data.value ?? []) {
ensureNotAborted(options.signal, "loadOneDrive");
if (docs.length >= maxFiles) return;
if (item.folder) {
await walk(`items/${item.id}`);
continue;
}
if (!item.file) continue;
if (options.filter && !options.filter(item.name)) continue;
const downloadUrl = item["@microsoft.graph.downloadUrl"];
if (!downloadUrl) continue;
attempted++;
try {
const fileResponse = await doFetch(fetchImpl, downloadUrl, { signal: options.signal }, "loadOneDrive");
if (!fileResponse.ok) continue;
docs.push({
content: await readResponseText(fileResponse, "loadOneDrive"),
source: `onedrive:${item.id}`,
metadata: { id: item.id, name: item.name, mimeType: item.file.mimeType }
});
loaded++;
} catch (err) {
rethrowIfAbort(err, options.signal, "loadOneDrive");
}
}
const next = data["@odata.nextLink"];
if (!next) {
url = void 0;
continue;
}
if (!item.file) continue;
if (options.filter && !options.filter(item.name)) continue;
const downloadUrl = item["@microsoft.graph.downloadUrl"];
if (!downloadUrl) continue;
const fileResponse = await fetchImpl(downloadUrl);
if (!fileResponse.ok) continue;
docs.push({
content: await fileResponse.text(),
source: `onedrive:${item.id}`,
metadata: { id: item.id, name: item.name, mimeType: item.file.mimeType }
});
if (seenLinks.has(next) || next === url) {
throw loadFailed("loadOneDrive: incomplete pagination (repeated @odata.nextLink)");
}
url = next;
}
}
await walk(folder);
return docs;
return finishTreeLoad("loadOneDrive", attempted, loaded, docs);
}
async function loadPdf(url, options) {
const fetchImpl = options.fetch ?? globalThis.fetch;
const response = await fetchImpl(url);
if (!response.ok) throw new RagError({ code: RagErrorCodes.AK_RAG_LOAD_FAILED, message: `loadPdf ${response.status}: ${url}` });
const buf = new Uint8Array(await response.arrayBuffer());
const { text, pages } = await options.parsePdf(buf);
return [{ content: text, source: url, metadata: { url, pages } }];
}
// src/rerankers/voyage.ts
function rerankFailed(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function voyageReranker(options) {

@@ -324,23 +591,52 @@ const fetchImpl = options.fetch ?? globalThis.fetch;

if (documents.length === 0) return documents;
const response = await fetchImpl("https://api.voyageai.com/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
})
});
let response;
try {
response = await fetchImpl("https://api.voyageai.com/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
}),
signal: options.signal
});
} catch (cause) {
throw rerankFailed("voyage rerank: network error", cause);
}
if (!response.ok) {
const text = await response.text();
throw new RagError({ code: RagErrorCodes.AK_RAG_RERANK_FAILED, message: `voyage rerank: ${response.status} ${text.slice(0, 200)}` });
let text = "";
try {
text = await response.text();
} catch (cause) {
throw rerankFailed(`voyage rerank: ${response.status} (failed to read error body)`, cause);
}
throw rerankFailed(`voyage rerank: ${response.status} ${text.slice(0, 200)}`);
}
const data = await response.json();
const ranked = (data.data ?? []).sort((a, b) => b.relevance_score - a.relevance_score).map((r) => {
let data;
try {
data = await response.json();
} catch (cause) {
throw rerankFailed("voyage rerank: invalid JSON response", cause);
}
if (!Array.isArray(data.data)) {
throw rerankFailed("voyage rerank: data must be an array");
}
const ranked = [];
for (let i = 0; i < data.data.length; i++) {
const r = data.data[i];
if (r == null || typeof r !== "object" || !Number.isInteger(r.index) || r.index < 0 || r.index >= documents.length || typeof r.relevance_score !== "number" || !Number.isFinite(r.relevance_score)) {
throw rerankFailed(`voyage rerank: malformed result at index ${i}`);
}
const doc = documents[r.index];
return { ...doc, score: r.relevance_score };
});
ranked.push({
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0,
score: r.relevance_score
});
}
ranked.sort((a, b) => b.score - a.score);
return ranked;

@@ -351,2 +647,9 @@ };

// src/rerankers/jina.ts
function rerankFailed2(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function jinaReranker(options) {

@@ -357,23 +660,52 @@ const fetchImpl = options.fetch ?? globalThis.fetch;

if (documents.length === 0) return documents;
const response = await fetchImpl("https://api.jina.ai/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
})
});
let response;
try {
response = await fetchImpl("https://api.jina.ai/v1/rerank", {
method: "POST",
headers: {
"content-type": "application/json",
authorization: `Bearer ${options.apiKey}`
},
body: JSON.stringify({
query,
documents: documents.map((d) => d.content),
model
}),
signal: options.signal
});
} catch (cause) {
throw rerankFailed2("jina rerank: network error", cause);
}
if (!response.ok) {
const text = await response.text();
throw new RagError({ code: RagErrorCodes.AK_RAG_RERANK_FAILED, message: `jina rerank: ${response.status} ${text.slice(0, 200)}` });
let text = "";
try {
text = await response.text();
} catch (cause) {
throw rerankFailed2(`jina rerank: ${response.status} (failed to read error body)`, cause);
}
throw rerankFailed2(`jina rerank: ${response.status} ${text.slice(0, 200)}`);
}
const data = await response.json();
const ranked = (data.results ?? []).sort((a, b) => b.relevance_score - a.relevance_score).map((r) => {
let data;
try {
data = await response.json();
} catch (cause) {
throw rerankFailed2("jina rerank: invalid JSON response", cause);
}
if (!Array.isArray(data.results)) {
throw rerankFailed2("jina rerank: results must be an array");
}
const ranked = [];
for (let i = 0; i < data.results.length; i++) {
const r = data.results[i];
if (r == null || typeof r !== "object" || !Number.isInteger(r.index) || r.index < 0 || r.index >= documents.length || typeof r.relevance_score !== "number" || !Number.isFinite(r.relevance_score)) {
throw rerankFailed2(`jina rerank: malformed result at index ${i}`);
}
const doc = documents[r.index];
return { ...doc, score: r.relevance_score };
});
ranked.push({
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0,
score: r.relevance_score
});
}
ranked.sort((a, b) => b.score - a.score);
return ranked;

@@ -384,11 +716,85 @@ };

// src/rerank.ts
function resolvePositiveInt2(value, fallback) {
if (value === void 0 || !Number.isFinite(value)) return fallback;
return Math.max(1, Math.floor(value));
}
function resolveWeight(value, fallback) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return fallback;
return value;
}
function resolveRelativeWeights(vectorWeight, bm25Weight) {
const v = resolveWeight(vectorWeight, 0.6);
const b = resolveWeight(bm25Weight, 0.4);
if (v === 0 && b === 0) return { vectorWeight: 0.5, bm25Weight: 0.5 };
const m = Math.max(v, b);
const vN = v / m;
const bN = b / m;
const sum = vN + bN;
return { vectorWeight: vN / sum, bm25Weight: bN / sum };
}
function cloneDoc(doc) {
return {
...doc,
metadata: doc.metadata ? { ...doc.metadata } : void 0
};
}
function rerankFailed3(message, cause) {
return new RagError({
code: RagErrorCodes.AK_RAG_RERANK_FAILED,
message,
cause
});
}
function enforceScoreOrder2(docs) {
if (docs.length === 0) return docs;
const anyScore = docs.some((d) => d.score !== void 0);
if (!anyScore) return docs.map(cloneDoc);
for (const d of docs) {
if (typeof d.score !== "number" || !Number.isFinite(d.score)) {
throw rerankFailed3(
"reranker output: every document must have a finite numeric score when scores are present"
);
}
}
return docs.map(cloneDoc).sort((a, b) => b.score - a.score);
}
function validateRerankOutput(value) {
if (!Array.isArray(value)) {
throw rerankFailed3("reranker output must be an array of documents");
}
const out = [];
for (let i = 0; i < value.length; i++) {
const item = value[i];
if (item == null || typeof item !== "object") {
throw rerankFailed3(`reranker output[${i}] is not a document object`);
}
const doc = item;
if (typeof doc.id !== "string" || typeof doc.content !== "string") {
throw rerankFailed3(`reranker output[${i}] must have string id and content`);
}
if (doc.score !== void 0 && (typeof doc.score !== "number" || !Number.isFinite(doc.score))) {
throw rerankFailed3(`reranker output[${i}] has a non-finite score`);
}
out.push(cloneDoc(doc));
}
return out;
}
function createRerankedRetriever(base, options = {}) {
const candidatePool = Math.max(1, options.candidatePool ?? 20);
const topK = Math.max(1, options.topK ?? 5);
const candidatePool = resolvePositiveInt2(options.candidatePool, 20);
const topK = resolvePositiveInt2(options.topK, 5);
const rerank = options.rerank ?? bm25Rerank;
return {
async retrieve(request) {
const candidates = (await base.retrieve(request)).slice(0, candidatePool);
const reranked = await rerank({ query: request.query, documents: candidates });
return reranked.slice(0, topK);
const baseResults = await base.retrieve(request);
const candidates = baseResults.slice(0, candidatePool).map(cloneDoc);
let reranked;
try {
reranked = await rerank({ query: request.query, documents: candidates });
} catch (cause) {
if (cause instanceof RagError) throw cause;
throw rerankFailed3("reranker threw", cause);
}
const validated = validateRerankOutput(reranked);
const ordered = enforceScoreOrder2(validated);
return ordered.slice(0, topK);
}

@@ -400,8 +806,16 @@ };

}
function resolveBm25K1(value) {
if (value === void 0 || !Number.isFinite(value) || value < 0) return 1.5;
return value;
}
function resolveBm25B(value) {
if (value === void 0 || !Number.isFinite(value) || value < 0 || value > 1) return 0.75;
return value;
}
function bm25Score(query, documents, options = {}) {
const k1 = options.k1 ?? 1.5;
const b = options.b ?? 0.75;
const k1 = resolveBm25K1(options.k1);
const b = resolveBm25B(options.b);
const qTerms = tokenize(query);
const N = documents.length;
if (N === 0 || qTerms.length === 0) return documents;
if (N === 0 || qTerms.length === 0) return documents.map(cloneDoc);
const docTerms = documents.map((d) => tokenize(d.content));

@@ -426,5 +840,8 @@ const avgdl = docTerms.reduce((acc, t) => acc + t.length, 0) / N;

const norm = 1 - b + b * (dl / (avgdl || 1));
score += idf * (f * (k1 + 1) / (f + k1 * norm));
const denom = f + k1 * norm;
if (denom === 0) continue;
const term = idf * (f * (k1 + 1) / denom);
if (Number.isFinite(term)) score += term;
}
return { ...doc, score };
return { ...cloneDoc(doc), score: Number.isFinite(score) ? score : 0 };
});

@@ -435,18 +852,34 @@ return scored.sort((a, b2) => (b2.score ?? 0) - (a.score ?? 0));

function normalize(docs) {
const scores = docs.map((d) => d.score ?? 0);
const max = Math.max(...scores, 0);
const map = /* @__PURE__ */ new Map();
const entries = [];
for (const d of docs) {
map.set(d.id, max > 0 ? (d.score ?? 0) / max : 0);
const raw = d.score;
const score = typeof raw === "number" && Number.isFinite(raw) ? raw : 0;
entries.push({ id: d.id, score });
}
const scores = entries.map((e) => e.score);
const min = scores.length > 0 ? Math.min(...scores) : 0;
const max = scores.length > 0 ? Math.max(...scores) : 0;
const range = max - min;
const map = /* @__PURE__ */ new Map();
for (const e of entries) {
if (map.has(e.id)) continue;
if (range > 0) {
map.set(e.id, (e.score - min) / range);
} else {
map.set(e.id, max === 0 ? 0 : 1);
}
}
return map;
}
function createHybridRetriever(base, options = {}) {
const vectorWeight = options.vectorWeight ?? 0.6;
const bm25Weight = options.bm25Weight ?? 0.4;
const topK = Math.max(1, options.topK ?? 5);
const candidatePool = Math.max(1, options.candidatePool ?? 20);
const { vectorWeight, bm25Weight } = resolveRelativeWeights(
options.vectorWeight,
options.bm25Weight
);
const topK = resolvePositiveInt2(options.topK, 5);
const candidatePool = resolvePositiveInt2(options.candidatePool, 20);
return {
async retrieve(request) {
const candidates = (await base.retrieve(request)).slice(0, candidatePool);
const baseResults = await base.retrieve(request);
const candidates = baseResults.slice(0, candidatePool).map(cloneDoc);
if (candidates.length === 0) return candidates;

@@ -456,6 +889,9 @@ const vectorScores = normalize(candidates);

const bm25Scores = normalize(bm25Docs);
const merged = candidates.map((d) => ({
...d,
score: vectorWeight * (vectorScores.get(d.id) ?? 0) + bm25Weight * (bm25Scores.get(d.id) ?? 0)
}));
const merged = candidates.map((d) => {
const raw = vectorWeight * (vectorScores.get(d.id) ?? 0) + bm25Weight * (bm25Scores.get(d.id) ?? 0);
return {
...d,
score: Number.isFinite(raw) ? raw : 0
};
});
merged.sort((a, b) => (b.score ?? 0) - (a.score ?? 0));

@@ -462,0 +898,0 @@ return merged.slice(0, topK);

{
"name": "@agentskit/rag",
"version": "0.4.14",
"version": "0.4.16",
"description": "Plug-and-play retrieval-augmented generation for AgentsKit.",

@@ -51,8 +51,8 @@ "keywords": [

"dependencies": {
"@agentskit/core": "1.12.3"
"@agentskit/core": "1.12.4"
},
"devDependencies": {
"@types/node": "^25.9.5",
"metro": "0.84.4",
"metro-runtime": "0.84.4",
"metro": "0.87.0",
"metro-runtime": "0.87.0",
"tsup": "^8.5.1",

@@ -63,3 +63,4 @@ "typescript": "^6.0.3",

"publishConfig": {
"access": "public"
"access": "public",
"provenance": true
},

@@ -66,0 +67,0 @@ "agentskit": {

# @agentskit/rag
<p align="center"><img src="https://raw.githubusercontent.com/AgentsKit-io/agentskit/main/apps/docs-next/public/brand/logo-wordmark.svg" alt="AgentsKit" width="180" /></p>
Profile: <code>major-package</code>
<p align="center"><img alt="AgentsKit" src="https://raw.githubusercontent.com/AgentsKit-io/agentskit/main/apps/docs-next/public/brand/logo-wordmark.svg" width="180" /></p>
Plug-and-play retrieval-augmented generation: chunk documents, embed them, and retrieve the right context at query time.

@@ -16,2 +18,8 @@

## Verified proof
- Package metadata and tests live under `packages/rag/`.
- Package guide: https://www.agentskit.io/docs/packages/rag
- Stability map: [docs/STABILITY.md](../../docs/STABILITY.md)
## How this fits the ecosystem

@@ -37,2 +45,3 @@

<!-- readme-command:install -->
```bash

@@ -44,2 +53,3 @@ npm install @agentskit/rag @agentskit/memory @agentskit/adapters

<!-- readme-example:quickstart -->
```ts

@@ -60,2 +70,3 @@ import { createRAG } from '@agentskit/rag'

const docs = await rag.search('How does AgentsKit work?', { topK: 5 })
console.log(docs)
```

@@ -92,2 +103,5 @@

- **Document loaders:** `loadUrl`, `loadGitHubFile`, `loadGitHubTree`, `loadNotionPage`, `loadConfluencePage`, `loadGoogleDriveFile`, `loadPdf` (BYO parser). [Recipe](https://www.agentskit.io/docs/recipes/doc-loaders).
- **Loader resilience:** HTTP/network and response-body read/parse failures surface as `RagError` (`AK_RAG_LOAD_FAILED`). Optional `signal` aborts (including mid-body reads) are never swallowed as a per-object skip. Tree/list loaders may return partial success when at least one eligible download succeeded; if every attempted eligible download failed, they throw. Missing/invalid S3 object bodies count as failed downloads. Pagination that reports more data without a new cursor/token throws (no silent truncation). `loadNotionPage` follows Notion `has_more` / `next_cursor` with `start_cursor` until complete (preserving block order; incomplete or repeated cursors throw). Non-positive / non-finite `maxFiles` yields `[]`.
- **Score contracts:** scoreless search/rerank results keep order. When any score is present, every result must have a finite numeric score and is sorted descending — mixed or non-finite scores throw (never fabricate `-Infinity`). Malformed Voyage/Jina/custom reranker output throws `AK_RAG_RERANK_FAILED`. Optional `signal` on `voyageReranker` / `jinaReranker` is forwarded to `fetch`; request/body aborts remain `AK_RAG_RERANK_FAILED`. `bm25Score` sanitizes invalid `k1`/`b` to documented defaults and always emits finite scores. Hybrid relative weights are normalized to a finite pair that sums to 1 (both zero → 0.5/0.5).
- **Chunk/config safety:** invalid `chunkSize` / `chunkOverlap` / `topK` values are sanitized so chunking always terminates and search never sends non-finite limits to the store.

@@ -132,1 +146,11 @@ ### S3 in Expo and React Native runtimes

[Full documentation](https://www.agentskit.io) · [GitHub](https://github.com/AgentsKit-io/agentskit)
## Maturity and compatibility
- Stability: **beta** — see [docs/STABILITY.md](../../docs/STABILITY.md)
- **Node.js 20+** and **TypeScript** strict mode
- Published as `@agentskit/rag`
## Contributing
See [CONTRIBUTING.md](../../CONTRIBUTING.md) and the monorepo [LICENSE](../../LICENSE).
// src/chunker.ts
function chunkText(text, options) {
if (!text) return [];
if (options.split) {
return options.split(text).filter((chunk) => chunk.length > 0);
}
const { chunkSize, chunkOverlap } = options;
if (text.length <= chunkSize) return [text];
const chunks = [];
let start = 0;
while (start < text.length) {
let end = Math.min(start + chunkSize, text.length);
if (end < text.length) {
const boundary = text.lastIndexOf(" ", end);
if (boundary > start) {
end = boundary;
}
}
const chunk = text.slice(start, end).trim();
if (chunk.length > 0) {
chunks.push(chunk);
}
if (end >= text.length) break;
const advance = end - start - chunkOverlap;
start += Math.max(advance, 1);
}
return chunks;
}
export { chunkText };
//# sourceMappingURL=chunk-UNFVK5RA.js.map
//# sourceMappingURL=chunk-UNFVK5RA.js.map
{"version":3,"sources":["../src/chunker.ts"],"names":[],"mappings":";AAMO,SAAS,SAAA,CAAU,MAAc,OAAA,EAAiC;AACvE,EAAA,IAAI,CAAC,IAAA,EAAM,OAAO,EAAC;AAEnB,EAAA,IAAI,QAAQ,KAAA,EAAO;AACjB,IAAA,OAAO,OAAA,CAAQ,MAAM,IAAI,CAAA,CAAE,OAAO,CAAA,KAAA,KAAS,KAAA,CAAM,SAAS,CAAC,CAAA;AAAA,EAC7D;AAEA,EAAA,MAAM,EAAE,SAAA,EAAW,YAAA,EAAa,GAAI,OAAA;AAEpC,EAAA,IAAI,IAAA,CAAK,MAAA,IAAU,SAAA,EAAW,OAAO,CAAC,IAAI,CAAA;AAE1C,EAAA,MAAM,SAAmB,EAAC;AAC1B,EAAA,IAAI,KAAA,GAAQ,CAAA;AAEZ,EAAA,OAAO,KAAA,GAAQ,KAAK,MAAA,EAAQ;AAC1B,IAAA,IAAI,MAAM,IAAA,CAAK,GAAA,CAAI,KAAA,GAAQ,SAAA,EAAW,KAAK,MAAM,CAAA;AAEjD,IAAA,IAAI,GAAA,GAAM,KAAK,MAAA,EAAQ;AACrB,MAAA,MAAM,QAAA,GAAW,IAAA,CAAK,WAAA,CAAY,GAAA,EAAK,GAAG,CAAA;AAC1C,MAAA,IAAI,WAAW,KAAA,EAAO;AACpB,QAAA,GAAA,GAAM,QAAA;AAAA,MACR;AAAA,IACF;AAEA,IAAA,MAAM,QAAQ,IAAA,CAAK,KAAA,CAAM,KAAA,EAAO,GAAG,EAAE,IAAA,EAAK;AAC1C,IAAA,IAAI,KAAA,CAAM,SAAS,CAAA,EAAG;AACpB,MAAA,MAAA,CAAO,KAAK,KAAK,CAAA;AAAA,IACnB;AAEA,IAAA,IAAI,GAAA,IAAO,KAAK,MAAA,EAAQ;AAExB,IAAA,MAAM,OAAA,GAAU,MAAM,KAAA,GAAQ,YAAA;AAC9B,IAAA,KAAA,IAAS,IAAA,CAAK,GAAA,CAAI,OAAA,EAAS,CAAC,CAAA;AAAA,EAC9B;AAEA,EAAA,OAAO,MAAA;AACT","file":"chunk-UNFVK5RA.js","sourcesContent":["export interface ChunkOptions {\n chunkSize: number\n chunkOverlap: number\n split?: (text: string) => string[]\n}\n\nexport function chunkText(text: string, options: ChunkOptions): string[] {\n if (!text) return []\n\n if (options.split) {\n return options.split(text).filter(chunk => chunk.length > 0)\n }\n\n const { chunkSize, chunkOverlap } = options\n\n if (text.length <= chunkSize) return [text]\n\n const chunks: string[] = []\n let start = 0\n\n while (start < text.length) {\n let end = Math.min(start + chunkSize, text.length)\n\n if (end < text.length) {\n const boundary = text.lastIndexOf(' ', end)\n if (boundary > start) {\n end = boundary\n }\n }\n\n const chunk = text.slice(start, end).trim()\n if (chunk.length > 0) {\n chunks.push(chunk)\n }\n\n if (end >= text.length) break\n\n const advance = end - start - chunkOverlap\n start += Math.max(advance, 1)\n }\n\n return chunks\n}\n"]}

Sorry, the diff of this file is too big to display

Sorry, the diff of this file is too big to display

Sorry, the diff of this file is too big to display

Sorry, the diff of this file is too big to display