RETRIEVAL 检索链路与数据库适配
English · SEARCH 契约与部署 · MEMORY
这些实现都在可选入口 @codesoul-co/ditto-retrieval 中。只冻结 RETRIEVAL.SEARCH 一个 Node;embedding、融合、重排是可组合的内部 Provider,也可以直接供 MEMORY 的进程内搜索使用。没有新增数据库驱动、数据库实例或索引管理模块。
存储与检索的接线
同一个应用插件可以实现 MemoryStore 和 MemorySearchProvider,两者使用同一个 SQL/Milvus 客户端。MEMORY 的 GET/QUERY/WRITE/UPDATE/DELETE 仍通过 store 操作原数据库。SEARCH 可以有三种接线:
- 直接注入已有数据库原生 MemorySearchProvider。
- 用
createRetrievalMemorySearchProvider将本文的检索链路适配为 MemorySearchProvider,无需注册 RETRIEVAL Worker。 - 用
RemoteRetrievalSearchProvider委托给已注册的 RETRIEVAL.SEARCH,可通过现有 HTTP transport 跨进程执行。
数据库已经支持文本到向量时,使用 nativeEmbedding: true 把文本交给数据库。否则,外部 EmbeddingProvider 在查询时生成向量,再交给数据库执行近邻搜索。查询向量与已存向量必须使用一致的模型、版本、维度及预处理约定;维度校验只能发现长度错误,不能识别两个等维但不兼容的模型。
批量构建文档向量可以显式调用 embedContents(..., { purpose: "document" }),然后由应用通过存储插件写入。SEARCH 不会自动写回向量、重建索引或同步 SQL 与 Milvus 数据。
Provider API
interface RetrievalExecutionContext {
signal?: AbortSignal;
defaults?: RetrievalDefaults;
}
interface EmbeddingProvider {
embed(input: {
contents: readonly unknown[];
purpose: "query" | "document";
}, context?: RetrievalExecutionContext): Promise<readonly (readonly number[])[]>;
}
interface RerankProvider {
rank(input: {
query: RetrievalQuery;
candidates: readonly RetrievalCandidate[];
limit: number;
}, context?: RetrievalExecutionContext): Promise<readonly {
index: number; score?: number;
}[]>;
}RetrievalSearchProvider.search(input, context?) 的第二参数可选,已有单参数实现继续兼容。Provider 返回原始结果;仅 Worker/SDK 包装 NodeResult。取消信号是进程内协作式取消,数据库回调需要自行传给驱动;Runtime 的 HTTP 调用不会自动把它传到远端。
| 导出 | 参数与行为 |
|---|---|
embedContents(provider, input, options?, context?) | options 为 { batchSize?, dimensions? };顺序分批、保持输入顺序;空输入不调用模型;校验数量、有限数值、非空向量和跨批次维度。dimensions 是期望输出维度,不会自动截断向量或修改模型请求。 |
createVectorSearchProvider({ backend, embedding?, nativeEmbedding?, dimensions? }) | 已有数值向量直接搜索;普通内容先 embedding,或在 nativeEmbedding 模式直接给后端。外部与原生 embedding 不能同时开启。 |
createTextSearchProvider(backend) | 校验非空字符串,全文/BM25 在数据库或搜索引擎执行;不做 embedding。 |
createHybridSearchProvider({ branches, candidateLimit?, rrfK?, key? }) | branches 为 { provider, strategy, weight? }[],并行搜索后执行加权 RRF;weight 默认 1,必须大于 0。 |
rerankCandidates(provider, input, context?) | rank 返回原候选的零起始 index;校验唯一性、范围、结果数量与分数。不会接受重排器生成的新内容。 |
createRerankSearchProvider({ search, reranker, candidateLimit? }) | 扩大召回池,重排后截取请求 limit。 |
createCosineReranker(embedding) | 对 query.content 和候选 content 分别 embedding,按余弦相似度降序排序;分数相同保持原顺序,零向量分数为 0。需要更复杂的 cross-encoder 时替换 RerankProvider。 |
HTTP embedding 接口只接受非空字符串;通用 EmbeddingProvider 允许结构化内容。若候选 content 是完整 MemoryItem,传给 cosine reranker 的 embedding 适配器需要显式提取 memory.content 文本,不能把整个对象直接交给 HTTP Provider。
融合分数为 sum(weight / (rrfK + rank)),rank 从 1 开始。不同检索器的原始分数不直接相加。默认身份为 [source.target ?? target.name, id ?? source.ref];同一分支的重复身份只计一次,使用首次出现的原始排名。缺少 id/ref 时必须提供 key。不同来源若代表同一记录,应由 key 统一身份。等分时保持首次遇到的顺序。
融合候选沿用首次出现的 content/source/metadata,score 替换为融合分数;每一分支的原始 score、rank、weight、source、metadata 留在 output.metadata.fusion.contributions 中,键为融合身份。分支输出 metadata 留在 fusion.branches 中。任一分支失败,整个搜索失败,并等待其他已开始的分支结束,不返回未声明的部分成功。
云端与本地 Provider
Provider 按能力区分为 EmbeddingProvider、RetrievalSearchProvider 和 RerankProvider,不再按部署位置增加平行的 Cloud/Local 类层次:
- 云端 HTTP 与本地 HTTP:复用 createHttpEmbeddingProvider,更换 baseUrl/model/apiKey;本地无认证服务可省略 apiKey,仍遵守网络权限配置。
- 云端 SDK 与进程内模型:直接实现 embed/search/rank,内部调用已有 SDK;模型加载、连接池、设备选择、批处理和关闭由应用拥有的实例负责,不逐请求创建模型或客户端。
- 数据库原生 embedding、全文或图检索:直接调用原 SDK 能力;只有后端需要向量时才在前面添加 embedding。
- 独立搜索引擎或远程检索 API:直接实现 RetrievalSearchProvider,按实际需要选择是否添加融合/重排,不强制走向量流程。
// encodeBatch 是应用已加载的本地模型或云端 SDK 的适配函数;没有额外 Worker。
const embedding: EmbeddingProvider = {
embed: ({ contents, purpose }, context) =>
encodeBatch(contents, { purpose, signal: context?.signal }),
};
const provider = createVectorSearchProvider({ backend: databaseSearch, embedding });
// 普通部署:在 MEMORY 内使用相同 Provider。
const memory = createMemoryWorker({
store: databaseStore,
search: createRetrievalMemorySearchProvider({
provider, target: { name: "agent-memory" }, defaults: config.retrieval,
}),
});
// 需要独立资源时:把 provider 注册到 RETRIEVAL registry,
// MEMORY 改为 search: new RemoteRetrievalSearchProvider({ runtime, target })。这里改变的是执行位置,数据库 SDK、embedding 模型和结果映射可保持不变。注册到独立进程时,由该进程创建自己的 SDK/模型资源;不要传输连接对象。自动负载检测和自动启动 Worker 不在 Provider 内实现。
HTTP embedding 与根目录配置
import { Sandbox } from "@codesoul-co/ditto/runtime/sandbox";
import { createHttpEmbeddingProvider, embeddingConfigFromEnv } from "@codesoul-co/ditto-retrieval";
const embedding = createHttpEmbeddingProvider({
...embeddingConfigFromEnv(process.env),
sandbox: new Sandbox(config.workspace, config.sandbox),
timeoutMs: config.timeoutMs,
});HttpEmbeddingOptions 为 { baseUrl, model, apiKey?, sandbox, timeoutMs?, fetch? },timeoutMs 默认 30000。也可注入自身的 Pick<Sandbox, "assert"> 权限对象。请求前执行 sandbox.assert("network", origin);禁止重定向。使用 OpenAI-compatible POST <baseUrl>/embeddings,请求 { model, input: string[], encoding_format: "float" };按响应 data.index 恢复顺序并校验向量。供应商其他协议可实现 EmbeddingProvider。
.env.example 的 DITTO_WORKER_RETRIEVAL_EMBEDDING_BASE_URL/API_KEY/MODEL 仅由显式调用 embeddingConfigFromEnv 消费。URL 和 model 必填,API key 可为空(例如本地服务);网络 origin 还需要加入共享 Sandbox 白名单。数据库地址、凭据和 SDK 初始化继续由应用插件管理。
workers:
retrieval:
searchLimit: 10
embedding:
batchSize: 64
# dimensions: 1536 # 可选,必须与索引和模型真实输出一致
hybrid:
candidateLimit: 100
rrfK: 60
rerank:
candidateLimit: 100配置加载为 config.retrieval。limit 优先级是请求 > Worker defaults > YAML > 10。其他数值优先级是 Provider 构造参数 > Worker defaults 对应字段 > YAML > 内置值。两个 candidateLimit 都是召回池下限,实际至少为请求 limit;嵌套重排时,重排需要的数量会传入内部搜索。
batchSize 范围 1..2048;limit/candidateLimit 范围 1..10000;dimensions/rrfK 必须为正安全整数。数值策略放 YAML;超时复用 runtime.timeoutMs,并显式传给 HTTP Provider。直接调用组合 Provider 时传 { defaults: config.retrieval } 作为 context;直接 SDK 则用 createRetrieval({ providers, defaults: config.retrieval })。
SQL 适配:MySQL / PostgreSQL
createSqlSearchProvider<Row>({ prepare, query, mapRow }) 不持有连接。prepare 接收完整 RetrievalSearchInput,返回 { text: string, values: readonly unknown[] };query 接收该语句和 context,返回 rows;mapRow 将每行转换为 RetrievalCandidate。所有返回结果仍受 target/limit/候选结构校验。
// 使用应用已有连接池,二选一;这些库不在 Ditto 的依赖中。
query: async ({ text, values }) => (await pgPool.query(text, [...values])).rows
query: async ({ text, values }) => (await mysqlPool.execute(text, [...values]))[0]应用按自己的表结构实现 prepare,绑定所有用户参数,明确过滤/租户权限、排序与 limit。例如表有 tenant/content 两列:
-- MySQL;values = [query, tenant, query, limit]
SELECT id, content, MATCH(content) AGAINST (? IN NATURAL LANGUAGE MODE) AS score
FROM memories WHERE tenant = ? AND MATCH(content) AGAINST (? IN NATURAL LANGUAGE MODE)
ORDER BY score DESC LIMIT ?
-- PostgreSQL;values = [query, tenant, limit]
SELECT id, content,
ts_rank_cd(to_tsvector('english', content), plainto_tsquery('english', $1)) AS score
FROM memories WHERE tenant = $2
AND to_tsvector('english', content) @@ plainto_tsquery('english', $1)
ORDER BY score DESC LIMIT $3这些示例需要应用已有的全文索引配置。MySQL 原生全文检索、PostgreSQL ts_rank 与 BM25 并非相同排名算法,应注册符合实际实现的 strategy,例如 keyword。需要 BM25 时接数据库真实 BM25 能力;测试使用 SQLite FTS5 的 bm25 函数验证此路径。
SQL 适配器不会猜测 schema 或把任意 filter 拼成 SQL。prepare 必须处理或拒绝所有传入的 filter/namespace。返回 score 应表达应用需要的相关性顺序,SQL 必须显式 ORDER BY;Ditto 保留数据库顺序。存储插件已实现原生 search 时,也可直接用下文的 createMemoryRetrievalProvider 复用,省去重复映射。
Milvus 适配
const database = createMilvusSearchProvider({
collection: "memories", vectorField: "embedding", outputFields: ["memory"],
search: (request, context) => applicationMilvusSearch(request, context),
// 应用已有客户端执行 client.search(request),应用封装负责驱动超时/取消。
scope(input) {
if (!input.target.namespace || input.filter) throw new Error("Unsupported scope");
return { filter: "tenant == {tenant}", exprValues: { tenant: input.target.namespace } };
},
mapHit: hit => ({
id: String(hit.id), content: hit.memory, score: hit.score,
source: { target: "agent-memory", ref: String(hit.id) },
}),
params: { ef: 64 },
});
const vector = createVectorSearchProvider({ backend: database, embedding });MilvusSearchOptions<Hit> 包含 collection/vectorField/outputFields/search/mapHit/scope?/params?。search 回调接收 { collection_name, anns_field, data, limit, output_fields, filter?, exprValues?, partition_names?, params? } 和 context;单次查询 data 为 [vector] 或 [text]。返回 { status: { code?, error_code?, reason? }, results: Hit[] },即 Milvus Node 客户端单查询展开结果;其它 SDK 响应由应用回调转换。非成功状态不会被误认为空召回。
应用负责 mapHit 的 ID、内容及 score 语义,例如 L2 distance 越小越相近,若需要“分数越高越好”应显式转换。该适配器不自行改分数。filter/namespace 出现时必须提供 scope mapper;不支持的条件应拒绝,不能静默忽略。options 不会被展开为 Milvus 请求,更不能覆盖固定 collection。
数据库已有文本 embedding 函数时,用 createVectorSearchProvider({ backend: database, nativeEmbedding: true })。数据库原生 BM25 也可用 createTextSearchProvider(database),将 vectorField 指向相应稀疏向量字段。具体模型函数、schema、索引和文本检索能力必须在应用数据库中预先配置;Ditto 不会替应用创建。
组合示例与 MEMORY 双向转接
const hybrid = createHybridSearchProvider({ branches: [
{ strategy: "vector", provider: vector },
{ strategy: "keyword", provider: createTextSearchProvider(sqlProvider) },
] });
const pipeline = createRerankSearchProvider({ search: hybrid, reranker: applicationReranker });
const providers = new RetrievalTargetRegistry({
"agent-memory": { defaultStrategy: "hybrid", providers: { vector, hybrid: pipeline } },
});
// 需要服务边界时注册:createRetrievalWorker({ providers })。
// Graph/temporal/custom 等策略直接注册自己的 SearchProvider,执行数据库遍历或专用算法;无需 embedding。MEMORY 转接函数来自 @codesoul-co/ditto-retrieval/adapters/memory,仅通过类型依赖 MEMORY:
| API | 用法 |
|---|---|
createMemoryRetrievalProvider(search, mapInput?) | 将已有 MemorySearchProvider 作为 RETRIEVAL 后端;结果 content 保留完整 MemoryItem。mapInput 可以适配 namespace 等请求信息;有 namespace 却未提供 mapper 会拒绝。 |
createRetrievalMemorySearchProvider({ provider, target, strategy?, defaults?, mapOutput? }) | 在 MEMORY 中直接运行检索 pipeline;defaults 显式传 config.retrieval。 |
RemoteRetrievalSearchProvider({ runtime, target, mapOutput? }) | 将同一 MEMORY.SEARCH 请求交给 Runtime 中的 RETRIEVAL.SEARCH;后端可以在另一个进程。 |
mapMemoryCandidates(output) | 默认结果映射:candidate.content 必须是完整 MemoryItem,candidate.id 若存在必须等于 memory.id;保留 key/content/metadata、score、命中 metadata 和 source。 |
若候选只携带 ID 或文本,自定义 mapOutput 一次批量补全/映射。默认映射不会从文本猜测 Memory ID。保持原生 search 插件方向为“数据库 → Retrieval Provider”,不能把已委托 RETRIEVAL 的 MEMORY.SEARCH 再作为同一 RETRIEVAL 的后端,避免递归。
逐 API 使用示例
完整代码:examples/retrieval.ts。下列函数共用该文件的 imports;函数不会在导入时自动执行。数据库、模型和 MCP 参数由应用注入,不是 Ditto 内置的模拟后端。选择需要的函数调用;写入、删除、模型调用等会产生对应的真实操作。
import { createDitto, createMemoryWorker, loadRuntimeConfigFile } from "@codesoul-co/ditto";
import type { MemorySearchProvider, MemoryStore, MemoryItem } from "@codesoul-co/ditto/worker/memory";
import {
createRetrieval, createRetrievalWorker, RetrievalTargetRegistry, RetrievalError, retrievalSearchNode,
embedContents, validateVector, createVectorSearchProvider, createTextSearchProvider,
createHybridSearchProvider, createCosineReranker, rerankCandidates, createRerankSearchProvider,
createHttpEmbeddingProvider, embeddingConfigFromEnv, createSqlSearchProvider, createMilvusSearchProvider,
type RetrievalSearchProvider, type RetrievalSearchInput, type RetrievalSearchOutput,
type EmbeddingProvider, type RerankProvider, type SqlSearchOptions, type MilvusSearchOptions,
} from "@codesoul-co/ditto-retrieval";
import {
createMemoryRetrievalProvider, createRetrievalMemorySearchProvider,
RemoteRetrievalSearchProvider, mapMemoryCandidates,
} from "@codesoul-co/ditto-retrieval/adapters/memory";
export const request: RetrievalSearchInput = { query: { content: "agent memory" }, target: { name: "kb" }, limit: 5 };EmbeddingProvider.embed / embedContents / validateVector
embed 是 SDK 适配接口,直接调用不附加批次包装;embedContents 顺序分批并校验维度。validateVector(value, dimensions?) 成功返回 void,失败抛 RETRIEVAL_INVALID_EMBEDDING。示例中的 dimensions 应等于实际模型输出维度。
export async function embeddingApis(embedding: EmbeddingProvider, dimensions: number) {
const direct = await embedding.embed({ contents: ["query text"], purpose: "query" });
const documents = await embedContents(embedding, { contents: ["document A", "document B"], purpose: "document" }, { batchSize: 32, dimensions });
validateVector(documents[0], dimensions); // Returns void; throws on invalid/empty/mismatched vectors.
return { direct, documents };
}embeddingConfigFromEnv / createHttpEmbeddingProvider
前者只读取 env 并返回 baseUrl/model/apiKey,后者创建实际 HTTP embedding Provider。函数需要 .env 已由 Node --env-file 或应用加载;Ditto 不隐式加载。YAML 仅保存 batchSize/dimensions 等行为参数。
export async function httpEmbedding() {
const config = loadRuntimeConfigFile("ditto.yaml", process.env);
const runtime = createDitto({ config });
try {
const embedding = createHttpEmbeddingProvider({ ...embeddingConfigFromEnv(process.env), sandbox: runtime.services.sandbox, timeoutMs: config.timeoutMs });
return await embedContents(embedding, { contents: ["query text"], purpose: "query" }, {}, { defaults: config.retrieval });
} finally { await runtime.close(); }
}createVectorSearchProvider:三种输入路径
三个函数分别展示外部 embedding、数据库原生 embedding、预计算向量。选择一种即可;后端能否接收文本取决于数据库配置。nativeEmbedding 与 embedding 不能同时设置。
export async function externalVector(backend: RetrievalSearchProvider, embedding: EmbeddingProvider) {
const vector = createVectorSearchProvider({ backend, embedding });
return vector.search(request);
}
export async function nativeVector(backend: RetrievalSearchProvider) {
return createVectorSearchProvider({ backend, nativeEmbedding: true }).search(request);
}
export async function precomputedVector(backend: RetrievalSearchProvider, vector: readonly number[]) {
return createVectorSearchProvider({ backend, dimensions: vector.length }).search({ ...request, query: { content: vector } });
}createTextSearchProvider:文本检索
返回 RetrievalSearchProvider。要求非空文本;不执行向量化,不自带 BM25/倒排索引算法。真正的检索、排序和过滤交给传入后端。
export async function textSearch(backend: RetrievalSearchProvider) {
return createTextSearchProvider(backend).search({ ...request, strategy: "keyword" });
}createHybridSearchProvider:融合
两个分支会并行运行,各自收到 branch.strategy 和扩大的 limit。成功返回经过 RRF 排序的原候选及 fusion 元数据。分支可包含不同数据库,但 key 必须能正确区分/统一候选身份。
export async function hybridSearch(vector: RetrievalSearchProvider, text: RetrievalSearchProvider) {
const hybrid = createHybridSearchProvider({ candidateLimit: 50, rrfK: 60,
branches: [{ provider: vector, strategy: "vector", weight: 1 }, { provider: text, strategy: "keyword", weight: 0.7 }],
});
return hybrid.search({ ...request, strategy: "hybrid" });
}RerankProvider.rank / rerankCandidates / createCosineReranker / createRerankSearchProvider
rank 输出 { index, score? }[];rerankCandidates 将索引还原为原候选。返回数量必须等于 min(limit, candidates.length),索引唯一且合法。cosine 使用候选 content 做 document embedding;若 content 是 MemoryItem,应在 embedding 适配器内提取实际文本。
export async function rerankApis(embedding: EmbeddingProvider, output: RetrievalSearchOutput) {
const reranker = createCosineReranker(embedding);
const input = { query: request.query, candidates: output.candidates, limit: Math.max(1, Math.min(3, output.candidates.length)) };
const order = await reranker.rank(input); // Raw zero-based indexes and scores.
const candidates = await rerankCandidates(reranker, input); // Original candidates in ranked order.
return { order, candidates };
}
export async function rerankSearch(search: RetrievalSearchProvider, reranker: RerankProvider) {
return createRerankSearchProvider({ search, reranker, candidateLimit: 50 }).search(request);
}createSqlSearchProvider:可注入 SQL SDK
此例假定 PostgreSQL memories(id, content, tenant) 表,并显式拒绝不支持的 filter/options。query 回调使用应用已有连接池。MySQL 使用同一工厂,但 prepare 必须换成该方言的占位符和查询;前文给出对应 SQL。本例按 namespace 绑定租户参数,权限校验仍由应用完成。
interface SqlRow { id: string; content: string; score: number; }
export function sqlSearch(query: SqlSearchOptions<SqlRow>["query"]) {
return createSqlSearchProvider<SqlRow>({
prepare(input) {
if (typeof input.query.content !== "string" || !input.target.namespace || input.filter || input.options) {
throw new RetrievalError("RETRIEVAL_INVALID_INPUT", "Expected text and namespace without additional filters or options");
}
return {
text: "SELECT id, content, ts_rank_cd(to_tsvector('english', content), plainto_tsquery('english', $1)) AS score FROM memories WHERE tenant = $2 AND to_tsvector('english', content) @@ plainto_tsquery('english', $1) ORDER BY score DESC, id ASC LIMIT $3",
values: [input.query.content, input.target.namespace, input.limit],
};
},
query, // e.g. async ({ text, values }) => (await pgPool.query(text, [...values])).rows
mapRow: row => ({ id: row.id, content: row.content, score: row.score }),
});
}createMilvusSearchProvider:可注入 Milvus SDK
将现有 SDK search 适配为固定请求/响应接口;status code=0 或 error_code=Success 才成功。mapHit 返回完整 MemoryItem,便于默认 MEMORY 转接。数据库内索引、字段、embedding 函数需预先存在。
interface MilvusHit { id: string | number; score: number; memory: MemoryItem; }
export function milvusSearch(search: MilvusSearchOptions<MilvusHit>["search"]) {
return createMilvusSearchProvider<MilvusHit>({ collection: "memories", vectorField: "embedding", outputFields: ["memory"],
search, // Adapt the application's SDK response to { status, results }.
scope(input) {
if (!input.target.namespace || input.filter || input.options) throw new RetrievalError("RETRIEVAL_INVALID_INPUT", "Expected a namespace without extra filters or options");
return { filter: "tenant == {tenant}", exprValues: { tenant: input.target.namespace } };
},
mapHit: hit => ({ id: String(hit.id), content: hit.memory, score: hit.score }),
});
}createMemoryRetrievalProvider / mapMemoryCandidates
原生 MemorySearchProvider → RetrievalSearchProvider → MemorySearchOutput。默认映射要求 candidate.content 是完整 MemoryItem;如果 candidate.id 存在,必须等于 content.id。
export async function nativeMemoryInRetrieval(nativeSearch: MemorySearchProvider) {
const provider = createMemoryRetrievalProvider(nativeSearch);
const output = await provider.search(request);
return mapMemoryCandidates(output); // content must be a complete MemoryItem, not just text.
}createMemoryRetrievalProvider 的 mapInput
有 target.namespace 时必须显式提供 mapInput,否则拒绝。下面只是一种应用 filter DSL,数据库插件必须明确支持 tenant 字段;不能把 namespace 当作已经鉴权的身份。
export function nativeMemoryWithNamespace(nativeSearch: MemorySearchProvider) {
return createMemoryRetrievalProvider(nativeSearch, input => {
if (!input.target.namespace) throw new RetrievalError("RETRIEVAL_INVALID_INPUT", "A namespace is required");
return {
query: input.query.content, filter: { ...input.filter, tenant: input.target.namespace },
...(input.limit === undefined ? {} : { limit: input.limit }),
...(input.strategy === undefined ? {} : { strategy: input.strategy }),
...(input.options === undefined ? {} : { options: input.options }),
};
}); // The native plugin must implement this tenant filter and enforce caller authorization.
}createRetrievalMemorySearchProvider:进程内复用
无需注册 RETRIEVAL Worker。默认 mapOutput 仍要求完整 MemoryItem 候选,provider 可以是 vector/text/hybrid/rerank 组合。策略优先级为 MEMORY 请求 strategy > 此构造参数 strategy;插件连接由应用拥有。
export async function localMemoryPipeline(store: MemoryStore, provider: RetrievalSearchProvider) {
const search = createRetrievalMemorySearchProvider({ provider, target: { name: "kb" }, strategy: "vector", defaults: { searchLimit: 5 } });
const runtime = createDitto({ workers: [createMemoryWorker({ store, search })] });
try { return await runtime.invoke("MEMORY.SEARCH", { query: "agent memory", limit: 3 }); }
finally { await runtime.close(); }
}RemoteRetrievalSearchProvider.search:委托执行
构造不打开连接、不启动 Worker。search 原始返回 MemorySearchOutput;放入 MEMORY Worker 后才包装 NodeResult。示例在同一 Runtime 内路由,换成 HTTP 远端注册不改 MEMORY 请求。
export async function delegatedMemory(store: MemoryStore, provider: RetrievalSearchProvider) {
const { runtime } = setupRetrieval(provider);
const search = new RemoteRetrievalSearchProvider({ runtime, target: { name: "kb" } });
runtime.register(createMemoryWorker({ store, search }));
try {
const raw = await search.search({ query: "agent memory", limit: 3 }); // Raw MemorySearchOutput.
const routed = await runtime.invoke("MEMORY.SEARCH", { query: "agent memory", limit: 3 });
return { raw, routed };
} finally { await runtime.close(); }
}mapOutput:自定义完整记录映射
当 candidate.content 只是业务内容时,可在 createRetrievalMemorySearchProvider 或 RemoteRetrievalSearchProvider 中传 mapOutput: mapTextCandidates。若索引只含片段或 ID,应改成一次批量读取真实完整记录;不要把片段冒充完整 Memory。
export function mapTextCandidates(output: RetrievalSearchOutput) {
return output.candidates.map(candidate => {
if (!candidate.id) throw new Error("The application requires candidate ids");
return { memory: { id: candidate.id, content: candidate.content }, ...(candidate.score === undefined ? {} : { score: candidate.score }) };
}); // Use as mapOutput only when content is the application's complete memory content.
}与 CONTEXT 组合
数据库检索、embedding、融合与重排 Provider 可注入 createRagStrategy 的 retrieve 步骤。CONTEXT 不要求单独启动 RETRIEVAL;检索后端生命周期仍由应用管理。 完整 CONTEXT API 与调用示例。
Context 检索适配器
从 @codesoul-co/ditto-retrieval/adapters/context 导入 createRetrievalContextStrategy 和 mapContextCandidates;Core 不导入这个可选模块。
createRetrievalContextStrategy(options) 返回 ContextRagStrategy。provider 与 runtime 不能同时设置:provider 在本地执行 RetrievalSearchProvider,可复用 SQL、Milvus、向量/embedding、hybrid 流水线;runtime 调用 RETRIEVAL.SEARCH;在 CONTEXT Worker 内可同时省略二者以继承本次执行的 Runtime,显式 runtime 始终优先。需预先注册本地或远端 RETRIEVAL Worker,按现有 direct/IPC/HTTP 方式部署。
| 参数 | 行为 |
|---|---|
| target | 必填逻辑 RetrievalTarget,构造时快照 |
| strategy | 可选检索策略名称 |
| defaults | RetrievalDefaults;用于本地流水线,也为两种模式提供缺省 limit;委托计算采用目标 Worker 配置 |
| mapInput(input) | 可选 ContextSelectInput → RetrievalSearchInput;显式处理 corpus、租户 namespace、filter 和授权 |
| mapOutput(output) | 可选完整 RetrievalSearchOutput → ContextItem[] 或 Promise;默认 mapContextCandidates |
默认映射要求 query,将 strategy.options 转为检索 options;limit 使用请求值、searchLimit 或 10,并限制最多 10000。显式提供 corpus 时必须配置 mapInput,避免丢失范围。limit=0 不调用后端;检索后仍由 Context 执行条数和 token 预算。
mapContextCandidates(output) 要求内容为 JSON。ID 由 target.name/type/namespace 及候选 ID 生成稳定哈希;无 ID 时用 source.ref,再回退 content。source.ref 映射为 source.uri,保留 metadata 和 score,补充 retrievalTarget。内容不合法抛 INVALID_PROVIDER_OUTPUT;业务记录需自定义 mapOutput。Context.SELECT 负责去重和预算,不覆盖缓存。例如 mapContextCandidates({ target: { name: "docs" }, candidates: [{ id: "1", content: "Evidence", score: 0.8 }] }) 返回一个稳定 ContextItem,含对应内容和 score 元数据。
import { createContext, type RuntimeClient } from "@codesoul-co/ditto";
import type { RetrievalSearchProvider } from "@codesoul-co/ditto-retrieval";
import { createRetrievalContextStrategy, mapContextCandidates } from "@codesoul-co/ditto-retrieval/adapters/context";
export function contextSearch(provider: RetrievalSearchProvider, runtime: RuntimeClient) {
const inline = createContext({ services: { ragStrategy: createRetrievalContextStrategy({
provider, target: { name: "documents" }, strategy: "vector",
}) } });
const delegated = createContext({ services: { ragStrategy: createRetrievalContextStrategy({
runtime, target: { name: "documents" }, strategy: "vector", mapOutput: mapContextCandidates,
}) } });
const input = { context: { items: [] }, purpose: "infer" as const, query: "agent runtime",
strategy: { kind: "rag" as const }, limit: 5 };
return Promise.all([inline.select(input), delegated.select(input)]);
}将策略注入 Context services.ragStrategy。每次调用的 signal 传入本地 Provider 或委托调用;未显式设置 runtime 时,内置 Worker 使用绑定当前执行的 Runtime,使已接收的 Graph 可以在关闭期间正常完成。网络取消目前停止调用方等待,不会终止远端数据库操作。可运行 SQLite FTS5 示例。
在 MEMORY Worker 内构造 RemoteRetrievalSearchProvider 时,可省略 runtime,继承本次执行绑定的 Runtime;显式 runtime 始终优先,独立 SDK 委托则必须提供它。例如 Worker 内可用 new RemoteRetrievalSearchProvider({ target: { name: "memories" } })。
