diff --git a/MemoryCore/README.md b/MemoryCore/README.md index 1d0e35b9b..82de9d37c 100644 --- a/MemoryCore/README.md +++ b/MemoryCore/README.md @@ -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. diff --git a/MemoryCore/openclaw.plugin.json b/MemoryCore/openclaw.plugin.json index c9af4d844..8fd0ded1a 100644 --- a/MemoryCore/openclaw.plugin.json +++ b/MemoryCore/openclaw.plugin.json @@ -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)" } } diff --git a/MemoryCore/src/config.ts b/MemoryCore/src/config.ts index 8702ea3a7..fe250d94d 100644 --- a/MemoryCore/src/config.ts +++ b/MemoryCore/src/config.ts @@ -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. */ @@ -608,6 +614,9 @@ export function parseConfig(raw: Record | 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, diff --git a/MemoryCore/src/core/hooks/auto-capture.ts b/MemoryCore/src/core/hooks/auto-capture.ts index 084ffb95b..1d224f241 100644 --- a/MemoryCore/src/core/hooks/auto-capture.ts +++ b/MemoryCore/src/core/hooks/auto-capture.ts @@ -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?.( diff --git a/MemoryCore/src/core/store/embedding.test.ts b/MemoryCore/src/core/store/embedding.test.ts new file mode 100644 index 000000000..83976133f --- /dev/null +++ b/MemoryCore/src/core/store/embedding.test.ts @@ -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((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 = {}): 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((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(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); + }); +}); diff --git a/MemoryCore/src/core/store/embedding.ts b/MemoryCore/src/core/store/embedding.ts index f07c34474..b71e07ba1 100644 --- a/MemoryCore/src/core/store/embedding.ts +++ b/MemoryCore/src/core/store/embedding.ts @@ -14,6 +14,7 @@ */ import type { Logger } from "../types.js"; +import { tiktokenCount } from "../../offload/context-token-tracker.js"; // ============================ // Types @@ -41,6 +42,12 @@ export interface OpenAIEmbeddingConfig { proxyUrl?: string; /** Max input text length in characters before truncation (default: 5000). */ maxInputChars?: number; + /** Token-aware input budget. Takes precedence over maxInputChars when set. */ + maxInputTokens?: number; + /** Retry count for 429, 5xx, network, and timeout failures (default: 2). */ + maxRetries?: number; + /** Initial exponential-backoff delay in milliseconds (default: 500). */ + retryBaseDelayMs?: number; /** Timeout per API call in milliseconds (default: 10000). */ timeoutMs?: number; } @@ -73,6 +80,11 @@ export interface EmbeddingService { embed(text: string, options?: EmbeddingCallOptions): Promise; /** Get embeddings for multiple texts (batched API call) */ embedBatch(texts: string[], options?: EmbeddingCallOptions): Promise; + /** Per-item result path for durable ingestion; successful items survive sibling failures. */ + embedBatchSettled?( + texts: string[], + options?: EmbeddingCallOptions, + ): Promise>; /** Return the configured vector dimensions */ getDimensions(): number; /** Return provider + model identifiers for change detection */ @@ -92,6 +104,18 @@ export interface EmbeddingService { startWarmup(): void; /** Optional: release resources (model memory, GPU, etc.) on shutdown */ close?(): void | Promise; + /** Safe operational state; never includes credentials or input text. */ + getHealth?(): EmbeddingHealth; +} + +export type EmbeddingFailureCategory = "rate_limit" | "server" | "timeout" | "network" | "response"; + +export interface EmbeddingHealth { + state: "ready" | "degraded" | "initializing"; + requests: number; + failures: number; + retries: number; + lastError?: { category: EmbeddingFailureCategory; at: string }; } /** @@ -365,8 +389,8 @@ export class LocalEmbeddingService implements EmbeddingService { /** Max texts per batch (OpenAI limit is 2048, we use a safe value) */ const MAX_BATCH_SIZE = 256; -/** Max retries for API calls */ -const MAX_RETRIES = 0; +/** Default retries for transient API calls. */ +const DEFAULT_MAX_RETRIES = 2; /** Default timeout per API call in milliseconds */ const DEFAULT_API_TIMEOUT_MS = 10_000; @@ -408,8 +432,12 @@ export class OpenAIEmbeddingService implements EmbeddingService { private readonly providerName: string; private readonly proxyUrl?: string; private readonly maxInputChars?: number; + private readonly maxInputTokens?: number; + private readonly maxRetries: number; + private readonly retryBaseDelayMs: number; private readonly timeoutMs: number; private readonly logger?: Logger; + private health: EmbeddingHealth = { state: "ready", requests: 0, failures: 0, retries: 0 }; constructor(config: OpenAIEmbeddingConfig, logger?: Logger) { if (!config.apiKey) { @@ -432,6 +460,15 @@ export class OpenAIEmbeddingService implements EmbeddingService { this.providerName = config.provider || "openai"; this.proxyUrl = config.proxyUrl?.trim() || undefined; this.maxInputChars = config.maxInputChars && config.maxInputChars > 0 ? config.maxInputChars : undefined; + this.maxInputTokens = config.maxInputTokens && config.maxInputTokens > 0 + ? Math.floor(config.maxInputTokens) + : undefined; + this.maxRetries = Number.isInteger(config.maxRetries) && config.maxRetries! >= 0 + ? Math.min(config.maxRetries!, 10) + : DEFAULT_MAX_RETRIES; + this.retryBaseDelayMs = config.retryBaseDelayMs && config.retryBaseDelayMs > 0 + ? Math.floor(config.retryBaseDelayMs) + : 500; this.timeoutMs = config.timeoutMs && config.timeoutMs > 0 ? config.timeoutMs : DEFAULT_API_TIMEOUT_MS; this.logger = logger; } @@ -454,6 +491,10 @@ export class OpenAIEmbeddingService implements EmbeddingService { // nothing to do — remote API is stateless } + getHealth(): EmbeddingHealth { + return { ...this.health, lastError: this.health.lastError ? { ...this.health.lastError } : undefined }; + } + async embed(text: string, options?: EmbeddingCallOptions): Promise { const [result] = await this.embedBatch([text], options); return result; @@ -463,9 +504,7 @@ export class OpenAIEmbeddingService implements EmbeddingService { if (texts.length === 0) return []; // Truncate texts exceeding maxInputChars limit - const processedTexts = this.maxInputChars - ? texts.map((t) => this.truncateInput(t)) - : texts; + const processedTexts = texts.map((text) => this.truncateInput(text)); // Split into sub-batches if needed if (processedTexts.length > MAX_BATCH_SIZE) { @@ -481,11 +520,46 @@ export class OpenAIEmbeddingService implements EmbeddingService { return this._callApi(processedTexts, options?.timeoutMs); } + async embedBatchSettled( + texts: string[], + options?: EmbeddingCallOptions, + ): Promise> { + const processedTexts = texts.map((text) => this.truncateInput(text)); + const results: Array<{ embedding?: Float32Array; error?: string }> = []; + for (const text of processedTexts) { + try { + const [embedding] = await this._callApi([text], options?.timeoutMs); + results.push({ embedding }); + } catch (error) { + results.push({ error: error instanceof Error ? error.message : String(error) }); + } + } + return results; + } + /** * Truncate input text to stay within the configured maxInputChars limit. * Logs a warning when truncation occurs. */ private truncateInput(text: string): string { + if (this.maxInputTokens) { + const originalTokens = tiktokenCount(text); + if (originalTokens <= this.maxInputTokens) return text; + const chars = Array.from(text); + let low = 0; + let high = chars.length; + while (low < high) { + const mid = Math.ceil((low + high) / 2); + if (tiktokenCount(chars.slice(0, mid).join("")) <= this.maxInputTokens) low = mid; + else high = mid - 1; + } + const truncated = chars.slice(0, low).join(""); + this.logger?.warn?.( + `${TAG} Input truncated from ${originalTokens} to ${tiktokenCount(truncated)} tokens ` + + `(maxInputTokens=${this.maxInputTokens}, chars=${text.length}->${truncated.length})`, + ); + return truncated; + } if (!this.maxInputChars || text.length <= this.maxInputChars) return text; this.logger?.warn?.( `${TAG} Input truncated from ${text.length} to ${this.maxInputChars} chars (maxInputChars limit)`, @@ -494,6 +568,7 @@ export class OpenAIEmbeddingService implements EmbeddingService { } private async _callApi(texts: string[], timeoutOverride?: number): Promise { + this.health.requests++; const body: Record = { input: texts, model: this.model, @@ -518,7 +593,7 @@ export class OpenAIEmbeddingService implements EmbeddingService { // Retry loop with timeout let lastError: Error | undefined; - for (let attempt = 0; attempt <= MAX_RETRIES; attempt++) { + for (let attempt = 0; attempt <= this.maxRetries; attempt++) { try { const controller = new AbortController(); const timeoutId = setTimeout(() => controller.abort(), timeoutOverride ?? this.timeoutMs); @@ -542,6 +617,10 @@ export class OpenAIEmbeddingService implements EmbeddingService { throw err; } lastError = err; + if (attempt < this.maxRetries) { + this.recordRetry(resp.status === 429 ? "rate_limit" : "server", attempt, resp.status); + await this.delayForRetry(attempt, resp.headers.get("retry-after")); + } continue; } @@ -553,27 +632,83 @@ export class OpenAIEmbeddingService implements EmbeddingService { // Sort by index to ensure correct order, then sanitize+normalize for consistency with local provider const sorted = [...json.data].sort((a, b) => a.index - b.index); - return sorted.map((d) => sanitizeAndNormalize(d.embedding)); + if (sorted.length !== texts.length) { + throw new Error(`Embedding API returned ${sorted.length} vectors for ${texts.length} inputs`); + } + const vectors = sorted.map((item, index) => { + if (!Array.isArray(item.embedding) || item.embedding.length !== this.dims) { + throw new Error( + `Embedding vector ${index} has dimension ${item.embedding?.length ?? "invalid"}; expected ${this.dims}`, + ); + } + if (item.embedding.some((value) => !Number.isFinite(value))) { + throw new Error(`Embedding vector ${index} contains non-finite values`); + } + const vector = sanitizeAndNormalize(item.embedding); + if (!vector.some((value) => value !== 0)) { + throw new Error(`Embedding vector ${index} has zero magnitude`); + } + return vector; + }); + this.health.state = "ready"; + return vectors; } finally { clearTimeout(timeoutId); } } catch (err) { // Non-retryable errors (4xx client errors) — rethrow immediately if (err instanceof EmbeddingApiError && err.isClientError()) { + this.recordFailure("response"); throw err; } lastError = err instanceof Error ? err : new Error(String(err)); // AbortError = timeout, retry - if (attempt < MAX_RETRIES) { - // Exponential backoff: 500ms, 1000ms - const delay = 500 * (attempt + 1); - await new Promise((r) => setTimeout(r, delay)); + if (attempt < this.maxRetries) { + const category = this.classifyError(lastError); + this.recordRetry(category, attempt); + await this.delayForRetry(attempt); } } } + this.recordFailure(this.classifyError(lastError)); throw lastError ?? new Error("Embedding API call failed after retries"); } + + private classifyError(error: Error | undefined): EmbeddingFailureCategory { + if (error instanceof EmbeddingApiError) { + return error.httpStatus === 429 ? "rate_limit" : "server"; + } + if (error?.name === "AbortError" || error?.name === "TimeoutError") return "timeout"; + if (error instanceof SyntaxError || error?.message.includes("vector")) return "response"; + return "network"; + } + + private recordRetry( + category: EmbeddingFailureCategory, + attempt: number, + status?: number, + ): void { + this.health.retries++; + this.logger?.warn?.( + `${TAG} transient embedding failure category=${category}${status ? ` status=${status}` : ""}; ` + + `retry=${attempt + 1}/${this.maxRetries}`, + ); + } + + private recordFailure(category: EmbeddingFailureCategory): void { + this.health.state = "degraded"; + this.health.failures++; + this.health.lastError = { category, at: new Date().toISOString() }; + } + + private async delayForRetry(attempt: number, retryAfter?: string | null): Promise { + const parsedRetryAfter = retryAfter && /^\d+$/.test(retryAfter) + ? Number(retryAfter) * 1000 + : undefined; + const delay = Math.min(parsedRetryAfter ?? this.retryBaseDelayMs * (2 ** attempt), 30_000); + await new Promise((resolve) => setTimeout(resolve, delay)); + } } // ============================ diff --git a/MemoryCore/src/core/store/factory.ts b/MemoryCore/src/core/store/factory.ts index 19b82e6f2..6d012efb1 100644 --- a/MemoryCore/src/core/store/factory.ts +++ b/MemoryCore/src/core/store/factory.ts @@ -107,6 +107,10 @@ export function createStoreBundle( dimensions: config.embedding.dimensions, sendDimensions: config.embedding.sendDimensions, maxInputChars: config.embedding.maxInputChars, + maxInputTokens: config.embedding.maxInputTokens, + maxRetries: config.embedding.maxRetries, + retryBaseDelayMs: config.embedding.retryBaseDelayMs, + timeoutMs: config.embedding.timeoutMs, }, logger); } diff --git a/MemoryCore/src/core/store/sqlite-vector-coverage.test.ts b/MemoryCore/src/core/store/sqlite-vector-coverage.test.ts new file mode 100644 index 000000000..a475d1b38 --- /dev/null +++ b/MemoryCore/src/core/store/sqlite-vector-coverage.test.ts @@ -0,0 +1,46 @@ +import fs from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { afterEach, describe, expect, it } from "vitest"; +import { VectorStore } from "./sqlite.js"; + +const tempDirs: string[] = []; + +afterEach(async () => { + await Promise.all(tempDirs.splice(0).map((dir) => fs.rm(dir, { recursive: true, force: true }))); +}); + +describe("SQLite vector coverage and replay", () => { + it("retains L0 source rows across failed embedding, restart, and idempotent backfill", async () => { + const dir = await fs.mkdtemp(path.join(os.tmpdir(), "tdai-vector-coverage-")); + tempDirs.push(dir); + const dbPath = path.join(dir, "vectors.db"); + let store = new VectorStore(dbPath, 3); + await store.init({ provider: "openai", model: "bge-m3" }); + await store.upsertL0({ + id: "l0-durable", + sessionKey: "session", + sessionId: "session", + role: "user", + messageText: "香港 🧬 code", + recordedAt: new Date(0).toISOString(), + timestamp: 0, + }); + expect(store.getVectorCoverage()).toEqual({ + l0: { sourceRows: 1, vectorRows: 0, coverage: 0 }, + l1: { sourceRows: 0, vectorRows: 0, coverage: 1 }, + }); + store.close(); + + store = new VectorStore(dbPath, 3); + await store.init({ provider: "openai", model: "bge-m3" }); + expect(store.getVectorCoverage().l0).toEqual({ sourceRows: 1, vectorRows: 0, coverage: 0 }); + + const vector = new Float32Array([1, 0, 0]); + await store.updateL0Embedding("l0-durable", vector); + await store.updateL0Embedding("l0-durable", vector); + expect(store.getVectorCoverage().l0).toEqual({ sourceRows: 1, vectorRows: 1, coverage: 1 }); + expect(await store.countL0()).toBe(1); + store.close(); + }); +}); diff --git a/MemoryCore/src/core/store/sqlite.ts b/MemoryCore/src/core/store/sqlite.ts index 20e9f4219..f71da8b9d 100644 --- a/MemoryCore/src/core/store/sqlite.ts +++ b/MemoryCore/src/core/store/sqlite.ts @@ -2291,6 +2291,31 @@ export class VectorStore implements IMemoryStore { } } + /** Report source persistence and vector-index completeness independently. */ + getVectorCoverage(): { + l0: { sourceRows: number; vectorRows: number; coverage: number }; + l1: { sourceRows: number; vectorRows: number; coverage: number }; + } { + const count = (table: string): number => { + try { + return Number((this.db.prepare(`SELECT COUNT(*) AS count FROM ${table}`).get() as { count: number }).count); + } catch { + return 0; + } + }; + const layer = (sourceRows: number, vectorRows: number) => ({ + sourceRows, + vectorRows, + coverage: sourceRows === 0 ? 1 : Math.min(1, vectorRows / sourceRows), + }); + const l0Source = count("l0_conversations"); + const l1Source = count("l1_records"); + return { + l0: layer(l0Source, this.vecTablesReady ? count("l0_vec") : 0), + l1: layer(l1Source, this.vecTablesReady ? count("l1_vec") : 0), + }; + } + /** * Re-embed all existing L1 and L0 texts with a new embedding function. * diff --git a/MemoryCore/src/core/store/store-pool.ts b/MemoryCore/src/core/store/store-pool.ts index 3279150cb..273f91ceb 100644 --- a/MemoryCore/src/core/store/store-pool.ts +++ b/MemoryCore/src/core/store/store-pool.ts @@ -380,7 +380,12 @@ export class StorePool { apiKey: embCfg.apiKey, model: embCfg.model, dimensions: embCfg.dimensions, + sendDimensions: embCfg.sendDimensions, maxInputChars: embCfg.maxInputChars, + maxInputTokens: embCfg.maxInputTokens, + maxRetries: embCfg.maxRetries, + retryBaseDelayMs: embCfg.retryBaseDelayMs, + timeoutMs: embCfg.timeoutMs, }, this.logger as StoreLogger); } diff --git a/MemoryCore/src/core/store/types.ts b/MemoryCore/src/core/store/types.ts index cc629c473..2a6a664db 100644 --- a/MemoryCore/src/core/store/types.ts +++ b/MemoryCore/src/core/store/types.ts @@ -565,6 +565,11 @@ export interface IMemoryStore extends MemoryPromptStore, MemoryGenerationRefStor init(providerInfo?: EmbeddingProviderInfo): MaybePromise; isDegraded(): boolean; getCapabilities(): StoreCapabilities; + /** Optional synchronous coverage snapshot for health endpoints. */ + getVectorCoverage?(): { + l0: { sourceRows: number; vectorRows: number; coverage: number }; + l1: { sourceRows: number; vectorRows: number; coverage: number }; + }; close(): void; // ── L1 Write ───────────────────────────────────────────── diff --git a/MemoryCore/src/gateway/config.ts b/MemoryCore/src/gateway/config.ts index 77eaa8e39..f825a1853 100644 --- a/MemoryCore/src/gateway/config.ts +++ b/MemoryCore/src/gateway/config.ts @@ -518,6 +518,9 @@ export function loadGatewayConfig(overrides?: Partial): GatewayCo sendDimensions: bool(topLevelEmbedding, "sendDimensions") ?? memory.embedding.sendDimensions, conflictRecallTopK: num(topLevelEmbedding, "conflictRecallTopK") ?? memory.embedding.conflictRecallTopK, maxInputChars: num(topLevelEmbedding, "maxInputChars") ?? memory.embedding.maxInputChars, + maxInputTokens: num(topLevelEmbedding, "maxInputTokens") ?? memory.embedding.maxInputTokens, + maxRetries: num(topLevelEmbedding, "maxRetries") ?? memory.embedding.maxRetries, + retryBaseDelayMs: num(topLevelEmbedding, "retryBaseDelayMs") ?? memory.embedding.retryBaseDelayMs, timeoutMs: num(topLevelEmbedding, "timeoutMs") ?? memory.embedding.timeoutMs, recallTimeoutMs: num(topLevelEmbedding, "recallTimeoutMs") ?? memory.embedding.recallTimeoutMs, captureTimeoutMs: num(topLevelEmbedding, "captureTimeoutMs") ?? memory.embedding.captureTimeoutMs, diff --git a/MemoryCore/src/gateway/server.ts b/MemoryCore/src/gateway/server.ts index 14dd041d9..ce5b47aaf 100644 --- a/MemoryCore/src/gateway/server.ts +++ b/MemoryCore/src/gateway/server.ts @@ -1373,13 +1373,17 @@ export class TdaiGateway { } private handleHealth(res: http.ServerResponse): void { + const vectorStore = this.core.getVectorStore(); + const embeddingService = this.core.getEmbeddingService(); const response: HealthResponse = { - status: this.core.getVectorStore() ? "ok" : "degraded", + status: vectorStore && embeddingService?.getHealth?.().state !== "degraded" ? "ok" : "degraded", version: VERSION, uptime: Math.floor((Date.now() - this.startTime) / 1000), stores: { - vectorStore: !!this.core.getVectorStore(), - embeddingService: !!this.core.getEmbeddingService(), + vectorStore: !!vectorStore, + embeddingService: !!embeddingService, + embeddingHealth: embeddingService?.getHealth?.(), + vectorCoverage: vectorStore?.getVectorCoverage?.(), }, // Integrated services status services: { diff --git a/MemoryCore/src/gateway/types.ts b/MemoryCore/src/gateway/types.ts index 21ff8ef11..fe9fd3384 100644 --- a/MemoryCore/src/gateway/types.ts +++ b/MemoryCore/src/gateway/types.ts @@ -2,6 +2,8 @@ * TDAI Gateway — Request/Response types for the HTTP API. */ +import type { EmbeddingHealth } from "../core/store/embedding.js"; + // ============================ // Common // ============================ @@ -22,6 +24,11 @@ export interface HealthResponse { stores: { vectorStore: boolean; embeddingService: boolean; + embeddingHealth?: EmbeddingHealth; + vectorCoverage?: { + l0: { sourceRows: number; vectorRows: number; coverage: number }; + l1: { sourceRows: number; vectorRows: number; coverage: number }; + }; }; /** Integrated services status (only present when state_backend is configured) */ services?: {