Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions MemoryCore/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,27 @@ Configuration templates:
- `tdai-gateway.yaml`: default Standalone + Skill configuration.
- `tdai-gateway.proxy.yaml`: LLM access through an OpenAI-compatible proxy.

For remote embeddings, prefer a token budget below the model's active context
instead of relying on character counts. Transient failures are retried with a
bounded exponential backoff, while permanent 4xx responses fail immediately:

```yaml
embedding:
provider: openai
baseUrl: http://127.0.0.1:11434/v1
apiKey: local
model: bge-m3
dimensions: 1024
sendDimensions: false
maxInputTokens: 7680 # leave margin below an 8192-token model context
maxRetries: 2
retryBaseDelayMs: 500
```

`/health` reports remote embedding failures and L0/L1 vector coverage separately
from durable source-row counts. Missing vectors can be replayed into existing
record IDs; source rows are not duplicated.

## Storage and isolation

- Memory and metadata are stored in SQLite.
Expand Down
7 changes: 5 additions & 2 deletions MemoryCore/openclaw.plugin.json
Original file line number Diff line number Diff line change
Expand Up @@ -105,8 +105,11 @@
"dimensions": { "type": "number", "description": "向量维度(必填,需与所选模型匹配)" },
"sendDimensions": { "type": "boolean", "default": true, "description": "是否在请求体中携带 dimensions 字段。默认 true(兼容 OpenAI text-embedding-3-* 的 Matryoshka 截断)。当目标服务为 BGE-M3 等不支持自定义维度的固定维度模型时,请设为 false,否则会被服务端以 HTTP 400 拒绝('does not support matryoshka representation')。" },
"conflictRecallTopK": { "type": "number", "default": 5, "description": "冲突检测时召回 Top-K 数" },
"maxInputChars": { "type": "number", "default": 5000, "description": "Embedding 输入文本最大字符数,超出时截断并打印警告日志(默认 5000,适合大多数模型的 token 上限)" },
"timeoutMs": { "type": "number", "default": 10000, "description": "单次 embedding API 调用超时(毫秒),超时后该次请求中止且不重试" },
"maxInputChars": { "type": "number", "default": 5000, "description": "Embedding 输入文本最大字符数(仅在未配置 maxInputTokens 时使用)" },
"maxInputTokens": { "type": "number", "description": "Embedding 输入 token 上限;优先于 maxInputChars,应低于模型上下文并预留特殊 token 余量" },
"maxRetries": { "type": "number", "default": 2, "description": "429、5xx、超时及网络错误的最大重试次数" },
"retryBaseDelayMs": { "type": "number", "default": 500, "description": "指数退避初始延迟(毫秒)" },
"timeoutMs": { "type": "number", "default": 10000, "description": "单次 embedding API 调用超时(毫秒)" },
"recallTimeoutMs": { "type": "number", "description": "recall 路径 embedding 超时(毫秒),覆盖 timeoutMs。用户等待中,建议设短一些(如 3000)" },
"captureTimeoutMs": { "type": "number", "description": "capture 路径 embedding 超时(毫秒),覆盖 timeoutMs。后台运行,可设长一些(如 15000)" }
}
Expand Down
9 changes: 9 additions & 0 deletions MemoryCore/src/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,12 @@ export interface EmbeddingConfig {
proxyUrl?: string;
/** Max input text length in characters before truncation (default: 5000). Texts exceeding this limit are truncated with a warning. */
maxInputChars: number;
/** Token-aware input budget. When set, takes precedence over maxInputChars. */
maxInputTokens?: number;
/** Retries for transient remote failures (default: 2). */
maxRetries: number;
/** Initial exponential retry delay in milliseconds (default: 500). */
retryBaseDelayMs: number;
/** Timeout per embedding API call in milliseconds (default: 10000). */
timeoutMs: number;
/** Override timeoutMs for recall-path embedding calls (user-facing, should be shorter). Falls back to timeoutMs. */
Expand Down Expand Up @@ -608,6 +614,9 @@ export function parseConfig(raw: Record<string, unknown> | undefined): MemoryTda
conflictRecallTopK: num(embeddingGroup, "conflictRecallTopK") ?? 5,
proxyUrl: embeddingProxyUrl,
maxInputChars: num(embeddingGroup, "maxInputChars") ?? 5000,
maxInputTokens: num(embeddingGroup, "maxInputTokens"),
maxRetries: num(embeddingGroup, "maxRetries") ?? 2,
retryBaseDelayMs: num(embeddingGroup, "retryBaseDelayMs") ?? 500,
timeoutMs: num(embeddingGroup, "timeoutMs") ?? 10_000,
recallTimeoutMs: num(embeddingGroup, "recallTimeoutMs") ?? undefined,
captureTimeoutMs: num(embeddingGroup, "captureTimeoutMs") ?? undefined,
Expand Down
14 changes: 12 additions & 2 deletions MemoryCore/src/core/hooks/auto-capture.ts
Original file line number Diff line number Diff line change
Expand Up @@ -259,12 +259,22 @@ export async function performAutoCapture(params: {
const tBgStart = performance.now();
try {
const texts = bgSnapshot.map((r) => r.content);
const embeddings = await bgEmbeddingService.embedBatch(texts);
const settled = bgEmbeddingService.embedBatchSettled
? await bgEmbeddingService.embedBatchSettled(texts)
: (await bgEmbeddingService.embedBatch(texts)).map((embedding) => ({ embedding }));

let bgUpdated = 0;
for (let i = 0; i < bgSnapshot.length; i++) {
const embedding = settled[i]?.embedding;
if (!embedding) {
bgLogger?.warn?.(
`${TAG} [L0-vec-index-bg] Embedding unavailable for ${bgSnapshot[i].recordId}; `
+ `source row retained for backfill (${settled[i]?.error ?? "unknown error"})`,
);
continue;
}
try {
const ok = await bgVectorStore.updateL0Embedding!(bgSnapshot[i].recordId, embeddings[i]);
const ok = await bgVectorStore.updateL0Embedding!(bgSnapshot[i].recordId, embedding);
if (ok) bgUpdated++;
} catch (err) {
bgLogger?.warn?.(
Expand Down
103 changes: 103 additions & 0 deletions MemoryCore/src/core/store/embedding.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
import http from "node:http";
import { afterEach, describe, expect, it } from "vitest";
import { tiktokenCount } from "../../offload/context-token-tracker.js";
import { OpenAIEmbeddingService } from "./embedding.js";

const servers: http.Server[] = [];

async function endpoint(
handler: (body: { input?: string[] }, attempt: number, res: http.ServerResponse) => void,
): Promise<{ baseUrl: string; attempts: () => number; inputs: string[][] }> {
let attempts = 0;
const inputs: string[][] = [];
const server = http.createServer((req, res) => {
const chunks: Buffer[] = [];
req.on("data", (chunk) => chunks.push(Buffer.from(chunk)));
req.on("end", () => {
const body = JSON.parse(Buffer.concat(chunks).toString("utf8")) as { input?: string[] };
inputs.push(body.input ?? []);
attempts++;
handler(body, attempts, res);
});
});
await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve));
servers.push(server);
const address = server.address();
if (!address || typeof address === "string") throw new Error("server did not bind");
return { baseUrl: `http://127.0.0.1:${address.port}/v1`, attempts: () => attempts, inputs };
}

function sendVectors(res: http.ServerResponse, count: number): void {
res.setHeader("content-type", "application/json");
res.end(JSON.stringify({
data: Array.from({ length: count }, (_, index) => ({ index, embedding: [1, 2, 3] })),
}));
}

function service(baseUrl: string, overrides: Record<string, unknown> = {}): OpenAIEmbeddingService {
return new OpenAIEmbeddingService({
provider: "openai",
baseUrl,
apiKey: "test",
model: "bge-m3",
dimensions: 3,
sendDimensions: false,
retryBaseDelayMs: 1,
...overrides,
});
}

afterEach(async () => {
await Promise.all(servers.splice(0).map((server) => new Promise<void>((resolve) => server.close(() => resolve()))));
});

describe("OpenAIEmbeddingService resilience", () => {
it("token-budgets mixed CJK, emoji, and code without splitting Unicode", async () => {
const remote = await endpoint((body, _attempt, res) => sendVectors(res, body.input?.length ?? 0));
const embedding = service(remote.baseUrl, { maxInputChars: 10_000, maxInputTokens: 32 });
const text = "香港資料庫🧬🚀 function example<T>(value: T) { return JSON.stringify(value); }".repeat(8);
await embedding.embed(text);
const sent = remote.inputs[0][0];
expect(tiktokenCount(sent)).toBeLessThanOrEqual(32);
expect(sent).not.toContain("�");
expect(sent.length).toBeLessThan(text.length);
});

it("retries bounded 500 and 429 responses with observable counters", async () => {
const remote = await endpoint((body, attempt, res) => {
if (attempt === 1) {
res.writeHead(500).end("server error");
} else if (attempt === 2) {
res.writeHead(429, { "retry-after": "0" }).end("rate limited");
} else sendVectors(res, body.input?.length ?? 0);
});
const embedding = service(remote.baseUrl, { maxRetries: 2 });
await expect(embedding.embed("retry me")).resolves.toHaveLength(3);
expect(remote.attempts()).toBe(3);
expect(embedding.getHealth()).toMatchObject({ state: "ready", requests: 1, failures: 0, retries: 2 });
});

it("retries timeouts only to the configured bound", async () => {
const remote = await endpoint((_body, _attempt, _res) => {
// Leave the response open until AbortSignal cancels the request.
});
const embedding = service(remote.baseUrl, { maxRetries: 1, timeoutMs: 15 });
await expect(embedding.embed("timeout")).rejects.toThrow();
expect(remote.attempts()).toBe(2);
expect(embedding.getHealth()).toMatchObject({ state: "degraded", requests: 1, failures: 1, retries: 1 });
expect(embedding.getHealth().lastError?.category).toBe("timeout");
});

it("preserves successful per-item progress and does not retry permanent 4xx", async () => {
const remote = await endpoint((body, _attempt, res) => {
if (body.input?.[0] === "bad") res.writeHead(400).end("invalid input");
else sendVectors(res, body.input?.length ?? 0);
});
const embedding = service(remote.baseUrl, { maxRetries: 3 });
const results = await embedding.embedBatchSettled(["bad", "good"]);
expect(results[0]).toMatchObject({ error: expect.stringContaining("HTTP 400") });
expect(results[1].embedding).toHaveLength(3);
expect(remote.attempts()).toBe(2);
expect(embedding.getHealth().failures).toBe(1);
});
});
Loading