fix(oss): v3 entity cleanup, filter fixes, and QA hardening (TS + Python) (#4858)
Co-authored-by: Soumil Rathi <soumilrathi@gmail.com> Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -262,6 +262,181 @@ export class Memory {
|
||||
return this._entityStore;
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalize a filters object for entity-store scoping: keeps only
|
||||
* user_id/agent_id/run_id keys whose values are defined.
|
||||
*/
|
||||
private _sessionFiltersFromPayload(
|
||||
payload: Record<string, any>,
|
||||
): Record<string, any> {
|
||||
const filters: Record<string, any> = {};
|
||||
if (payload.user_id) filters.user_id = payload.user_id;
|
||||
if (payload.agent_id) filters.agent_id = payload.agent_id;
|
||||
if (payload.run_id) filters.run_id = payload.run_id;
|
||||
return filters;
|
||||
}
|
||||
|
||||
/**
|
||||
* Remove `memoryId` from every entity record scoped to `filters`.
|
||||
* If an entity's `linkedMemoryIds` becomes empty after removal, the
|
||||
* entity record itself is deleted. Errors on individual entities are
|
||||
* swallowed so one bad record does not break the whole operation.
|
||||
*
|
||||
* No-op if the entity store has not been initialized yet.
|
||||
*/
|
||||
private async _removeMemoryFromEntityStore(
|
||||
memoryId: string,
|
||||
filters: Record<string, any>,
|
||||
): Promise<void> {
|
||||
let entityStore: VectorStore;
|
||||
try {
|
||||
entityStore = await this.getEntityStore();
|
||||
} catch (e) {
|
||||
console.debug(`Entity store unavailable during cleanup: ${e}`);
|
||||
return;
|
||||
}
|
||||
|
||||
let rows: Array<{ id: string; payload: Record<string, any> }> = [];
|
||||
try {
|
||||
const listed = await entityStore.list(filters, 10000);
|
||||
rows = (
|
||||
Array.isArray(listed) && Array.isArray(listed[0])
|
||||
? listed[0]
|
||||
: (listed as any)
|
||||
) as Array<{ id: string; payload: Record<string, any> }>;
|
||||
} catch (e) {
|
||||
console.debug(`Entity store list failed during cleanup: ${e}`);
|
||||
return;
|
||||
}
|
||||
|
||||
for (const row of rows) {
|
||||
try {
|
||||
const payload = row.payload || {};
|
||||
const linked: string[] = Array.isArray(payload.linkedMemoryIds)
|
||||
? payload.linkedMemoryIds
|
||||
: [];
|
||||
if (!linked.includes(memoryId)) continue;
|
||||
|
||||
const remaining = linked.filter((id) => id !== memoryId);
|
||||
if (remaining.length === 0) {
|
||||
try {
|
||||
await entityStore.delete(row.id);
|
||||
} catch (e) {
|
||||
console.debug(`Entity delete failed for id=${row.id}: ${e}`);
|
||||
}
|
||||
} else {
|
||||
const newPayload = { ...payload, linkedMemoryIds: remaining };
|
||||
// entityStore.update requires a vector — re-embed entity text.
|
||||
const entityText =
|
||||
typeof payload.data === "string" ? payload.data : "";
|
||||
if (!entityText) {
|
||||
// Can't re-embed without text; skip gracefully.
|
||||
console.debug(
|
||||
`Entity id=${row.id} missing 'data'; skipping update during cleanup`,
|
||||
);
|
||||
continue;
|
||||
}
|
||||
let vec: number[];
|
||||
try {
|
||||
vec = await this.embedder.embed(entityText);
|
||||
} catch (e) {
|
||||
console.debug(`Entity re-embed failed for '${entityText}': ${e}`);
|
||||
continue;
|
||||
}
|
||||
try {
|
||||
await entityStore.update(row.id, vec, newPayload);
|
||||
} catch (e) {
|
||||
console.debug(`Entity update failed for id=${row.id}: ${e}`);
|
||||
}
|
||||
}
|
||||
} catch (e) {
|
||||
console.debug(`Entity cleanup error for id=${row?.id}: ${e}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract entities from `text` and link them to `memoryId` in the
|
||||
* entity store, scoped to `filters` (user_id / agent_id / run_id).
|
||||
*
|
||||
* Simpler single-memory variant of Phase 7 in add(): no cross-memory
|
||||
* dedup, but still does per-entity "search for existing, update if
|
||||
* match >= 0.95 else insert new". Non-fatal errors are swallowed.
|
||||
*/
|
||||
private async _linkEntitiesForMemory(
|
||||
memoryId: string,
|
||||
text: string,
|
||||
filters: Record<string, any>,
|
||||
): Promise<void> {
|
||||
try {
|
||||
const entities = extractEntities(text);
|
||||
if (entities.length === 0) return;
|
||||
|
||||
const entityStore = await this.getEntityStore();
|
||||
|
||||
for (const entity of entities) {
|
||||
try {
|
||||
let entityVec: number[];
|
||||
try {
|
||||
entityVec = await this.embedder.embed(entity.text);
|
||||
} catch (e) {
|
||||
console.debug(`Entity embed failed for '${entity.text}': ${e}`);
|
||||
continue;
|
||||
}
|
||||
|
||||
let matches: Array<{
|
||||
id: string;
|
||||
score?: number;
|
||||
payload: Record<string, any>;
|
||||
}> = [];
|
||||
try {
|
||||
matches = await entityStore.search(entityVec, 1, filters);
|
||||
} catch {}
|
||||
|
||||
if (matches.length > 0 && (matches[0].score ?? 0) >= 0.95) {
|
||||
const match = matches[0];
|
||||
const payload = match.payload || {};
|
||||
const linked = new Set<string>(
|
||||
Array.isArray(payload.linkedMemoryIds)
|
||||
? payload.linkedMemoryIds
|
||||
: [],
|
||||
);
|
||||
linked.add(memoryId);
|
||||
payload.linkedMemoryIds = Array.from(linked).sort();
|
||||
try {
|
||||
await entityStore.update(match.id, entityVec, payload);
|
||||
} catch (e) {
|
||||
console.debug(`Entity update failed for '${entity.text}': ${e}`);
|
||||
}
|
||||
} else {
|
||||
const entityPayload: Record<string, any> = {
|
||||
data: entity.text,
|
||||
entityType: entity.type,
|
||||
linkedMemoryIds: [memoryId],
|
||||
};
|
||||
if (filters.user_id) entityPayload.user_id = filters.user_id;
|
||||
if (filters.agent_id) entityPayload.agent_id = filters.agent_id;
|
||||
if (filters.run_id) entityPayload.run_id = filters.run_id;
|
||||
|
||||
try {
|
||||
await entityStore.insert(
|
||||
[entityVec],
|
||||
[uuidv4()],
|
||||
[entityPayload],
|
||||
);
|
||||
} catch (e) {
|
||||
console.debug(`Entity insert failed for '${entity.text}': ${e}`);
|
||||
}
|
||||
}
|
||||
} catch (e) {
|
||||
console.debug(`Entity link error for '${entity.text}': ${e}`);
|
||||
}
|
||||
}
|
||||
} catch (e) {
|
||||
console.warn(`Entity linking failed during update: ${e}`);
|
||||
}
|
||||
}
|
||||
|
||||
private buildSessionScope(filters: SearchFilters): string {
|
||||
const parts: string[] = [];
|
||||
for (const key of ["agent_id", "run_id", "user_id"].sort()) {
|
||||
@@ -863,17 +1038,23 @@ export class Memory {
|
||||
// Validate search parameters (before applying defaults)
|
||||
validateSearchParams(config.threshold, config.topK);
|
||||
|
||||
// Validate and trim entity IDs in filters
|
||||
const normalizedFilters = config.filters
|
||||
? {
|
||||
...config.filters,
|
||||
user_id: validateAndTrimEntityId(config.filters.user_id, "user_id"),
|
||||
agent_id: validateAndTrimEntityId(
|
||||
config.filters.agent_id,
|
||||
"agent_id",
|
||||
),
|
||||
run_id: validateAndTrimEntityId(config.filters.run_id, "run_id"),
|
||||
}
|
||||
// Validate and trim entity IDs in filters. Only include keys whose
|
||||
// validated value is defined — otherwise downstream vector stores
|
||||
// receive `agent_id: undefined` / `run_id: undefined` and fail
|
||||
// (Qdrant rejects the malformed match, pgvector binds NULL, Redis
|
||||
// emits a literal "undefined" string in TAG filters).
|
||||
const normalizedFilters: Record<string, any> = config.filters
|
||||
? Object.fromEntries(
|
||||
Object.entries({
|
||||
...config.filters,
|
||||
user_id: validateAndTrimEntityId(config.filters.user_id, "user_id"),
|
||||
agent_id: validateAndTrimEntityId(
|
||||
config.filters.agent_id,
|
||||
"agent_id",
|
||||
),
|
||||
run_id: validateAndTrimEntityId(config.filters.run_id, "run_id"),
|
||||
}).filter(([, v]) => v !== undefined),
|
||||
)
|
||||
: {};
|
||||
|
||||
await this._ensureInitialized();
|
||||
@@ -1194,13 +1375,17 @@ export class Memory {
|
||||
|
||||
const { topK = 20 } = config;
|
||||
|
||||
// Validate and trim entity IDs in filters
|
||||
const filters = {
|
||||
...(config.filters || {}),
|
||||
user_id: validateAndTrimEntityId(config.filters?.user_id, "user_id"),
|
||||
agent_id: validateAndTrimEntityId(config.filters?.agent_id, "agent_id"),
|
||||
run_id: validateAndTrimEntityId(config.filters?.run_id, "run_id"),
|
||||
};
|
||||
// Validate and trim entity IDs in filters. Drop keys that resolve to
|
||||
// undefined so downstream vector stores don't receive
|
||||
// `agent_id: undefined` / `run_id: undefined` and fail.
|
||||
const filters: Record<string, any> = Object.fromEntries(
|
||||
Object.entries({
|
||||
...(config.filters || {}),
|
||||
user_id: validateAndTrimEntityId(config.filters?.user_id, "user_id"),
|
||||
agent_id: validateAndTrimEntityId(config.filters?.agent_id, "agent_id"),
|
||||
run_id: validateAndTrimEntityId(config.filters?.run_id, "run_id"),
|
||||
}).filter(([, v]) => v !== undefined),
|
||||
);
|
||||
|
||||
await this._captureEvent("get_all", {
|
||||
topK,
|
||||
@@ -1260,6 +1445,7 @@ export class Memory {
|
||||
...metadata,
|
||||
data,
|
||||
hash: createHash("md5").update(data).digest("hex"),
|
||||
textLemmatized: lemmatizeForBm25(data),
|
||||
createdAt: new Date().toISOString(),
|
||||
};
|
||||
|
||||
@@ -1317,6 +1503,16 @@ export class Memory {
|
||||
newMetadata.updatedAt,
|
||||
);
|
||||
|
||||
// Entity-store cleanup: strip this memory's id from old-text entities,
|
||||
// then re-extract entities from the new text and link them back.
|
||||
try {
|
||||
const sessionFilters = this._sessionFiltersFromPayload(newMetadata);
|
||||
await this._removeMemoryFromEntityStore(memoryId, sessionFilters);
|
||||
await this._linkEntitiesForMemory(memoryId, data, sessionFilters);
|
||||
} catch (e) {
|
||||
console.warn(`Entity store cleanup/link failed during update: ${e}`);
|
||||
}
|
||||
|
||||
return memoryId;
|
||||
}
|
||||
|
||||
@@ -1327,6 +1523,9 @@ export class Memory {
|
||||
}
|
||||
|
||||
const prevValue = existingMemory.payload.data;
|
||||
const sessionFilters = this._sessionFiltersFromPayload(
|
||||
existingMemory.payload || {},
|
||||
);
|
||||
await this.vectorStore.delete(memoryId);
|
||||
await this.db.addHistory(
|
||||
memoryId,
|
||||
@@ -1338,6 +1537,14 @@ export class Memory {
|
||||
1,
|
||||
);
|
||||
|
||||
// Entity-store cleanup: strip this memory's id from any entity records
|
||||
// that linked to it. Non-fatal — log and continue on error.
|
||||
try {
|
||||
await this._removeMemoryFromEntityStore(memoryId, sessionFilters);
|
||||
} catch (e) {
|
||||
console.warn(`Entity store cleanup failed during delete: ${e}`);
|
||||
}
|
||||
|
||||
return memoryId;
|
||||
}
|
||||
|
||||
|
||||
@@ -660,6 +660,9 @@ export function extractEntities(text: string): ExtractedEntity[] {
|
||||
txt = txt.replace(/\s*:+$/, "");
|
||||
// Strip leading numbered list markers
|
||||
txt = txt.replace(/^\d+\s*\.\s*/, "");
|
||||
// Strip trailing sentence punctuation (".", ",", ";", "!", "?") — otherwise
|
||||
// "Paris." and "Paris" produce different embeddings and break entity dedup.
|
||||
txt = txt.replace(/[.,;!?]+$/, "").trim();
|
||||
|
||||
if (!txt || txt.length <= 2 || hasArtifacts(txt)) {
|
||||
continue;
|
||||
|
||||
@@ -9,6 +9,19 @@ import type {
|
||||
import { VectorStore } from "./base";
|
||||
import { SearchFilters, VectorStoreConfig, VectorStoreResult } from "../types";
|
||||
|
||||
/**
|
||||
* Escape RediSearch TAG filter special characters. Any punctuation in the
|
||||
* value (including `-`, which appears in every UUID) must be backslash-
|
||||
* escaped, otherwise RediSearch either parses it as an operator (`-` is
|
||||
* minus, `|` is OR) or rejects the whole expression as a syntax error.
|
||||
*/
|
||||
function escapeRedisTagValue(value: unknown): string {
|
||||
return String(value).replace(
|
||||
/([,.<>{}\[\]"':;!@#$%^&*()\-+=~|/\\\s])/g,
|
||||
"\\$1",
|
||||
);
|
||||
}
|
||||
|
||||
interface RedisConfig extends VectorStoreConfig {
|
||||
redisUrl: string;
|
||||
collectionName: string;
|
||||
@@ -377,8 +390,8 @@ export class RedisDB implements VectorStore {
|
||||
const snakeFilters = filters ? toSnakeCase(filters) : undefined;
|
||||
const filterExpr = snakeFilters
|
||||
? Object.entries(snakeFilters)
|
||||
.filter(([_, value]) => value !== null)
|
||||
.map(([key, value]) => `@${key}:{${value}}`)
|
||||
.filter(([_, value]) => value !== null && value !== undefined)
|
||||
.map(([key, value]) => `@${key}:{${escapeRedisTagValue(value)}}`)
|
||||
.join(" ")
|
||||
: "*";
|
||||
|
||||
@@ -615,8 +628,8 @@ export class RedisDB implements VectorStore {
|
||||
const snakeFilters = filters ? toSnakeCase(filters) : undefined;
|
||||
const filterExpr = snakeFilters
|
||||
? Object.entries(snakeFilters)
|
||||
.filter(([_, value]) => value !== null)
|
||||
.map(([key, value]) => `@${key}:{${value}}`)
|
||||
.filter(([_, value]) => value !== null && value !== undefined)
|
||||
.map(([key, value]) => `@${key}:{${escapeRedisTagValue(value)}}`)
|
||||
.join(" ")
|
||||
: "*";
|
||||
|
||||
|
||||
Reference in New Issue
Block a user