Files
memind/memory-v2.mjs

435 lines
14 KiB
JavaScript

import { createMemoryV2PluginBackends } from './memory-v2-plugin-backends.mjs';
import { createPersonalMemoryShadowPipeline } from './memory-v2-personal-shadow.mjs';
const TRUE_VALUES = new Set(['1', 'true', 'yes', 'on']);
const FALSE_VALUES = new Set(['0', 'false', 'no', 'off']);
const MAX_MEMORY_TEXT_LENGTH = 2000;
const MAX_CONTEXT_ITEMS = 50;
function normalizeText(value, maxLength = MAX_MEMORY_TEXT_LENGTH) {
const text = String(value ?? '').trim();
if (!text) return null;
return text.length > maxLength ? text.slice(0, maxLength) : text;
}
function readFlag(env, key, fallback) {
const raw = env?.[key];
if (raw == null || raw === '') return fallback;
const normalized = String(raw).trim().toLowerCase();
if (TRUE_VALUES.has(normalized)) return true;
if (FALSE_VALUES.has(normalized)) return false;
return fallback;
}
function emptyResolveResult({
enabled,
skipped = false,
reason = null,
degraded = false,
source = null,
} = {}) {
return {
ok: !degraded,
enabled: Boolean(enabled),
skipped,
reason,
degraded,
source,
profile: null,
semanticMemories: [],
behaviorSummary: null,
activeGoals: [],
contextGoals: [],
memories: [],
};
}
function skippedWriteResult(reason) {
return {
ok: true,
enabled: false,
skipped: true,
reason,
source: null,
saved: 0,
analyzed: 0,
memories: 0,
};
}
export function resolveMemoryV2Policy({ env = process.env, overrides = {} } = {}) {
const legacyDefault = readFlag(env, 'USER_CONVERSATION_MEMORY_ENABLED', true);
const enabled = readFlag(env, 'MEMORY_ENABLED', legacyDefault);
return {
enabled,
profileEnabled: readFlag(env, 'MEMORY_PROFILE_ENABLED', enabled),
eventLogEnabled: readFlag(env, 'MEMORY_EVENT_LOG_ENABLED', enabled),
vectorEnabled: readFlag(env, 'MEMORY_VECTOR_ENABLED', false),
backend: String(env?.MEMORY_BACKEND ?? 'legacy').trim() || 'legacy',
failOpen: readFlag(env, 'MEMORY_FAIL_OPEN', true),
agentResolveEnabled: readFlag(env, 'MEMORY_AGENT_RESOLVE_ENABLED', false),
agentInjectionMode: normalizeAgentInjectionMode(env?.MEMORY_AGENT_INJECTION_MODE),
agentCanaryUserIds: normalizeUserIdList(env?.MEMORY_AGENT_CANARY_USER_IDS),
agentResolveLimit: resolveBoundedNumber(env?.MEMORY_AGENT_RESOLVE_LIMIT, 3, 1, 50),
agentResolveTimeoutMs: resolveBoundedNumber(env?.MEMORY_AGENT_RESOLVE_TIMEOUT_MS, 1200, 0, 30000),
promotionEnabled: readFlag(env, 'MEMORY_PROMOTION_ENABLED', false),
compactionV2Enabled: readFlag(env, 'MEMORY_COMPACTION_V2_ENABLED', false),
reflectionEnabled: readFlag(env, 'MEMORY_REFLECTION_ENABLED', false),
lifecycleWorkerEnabled: readFlag(env, 'MEMORY_LIFECYCLE_WORKER_ENABLED', false),
lifecycleRolloutMode: String(env?.MEMORY_LIFECYCLE_ROLLOUT_MODE ?? 'off').trim().toLowerCase() || 'off',
lifecycleRolloutUserIds: normalizeUserIdList(env?.MEMORY_LIFECYCLE_ROLLOUT_USER_IDS),
...overrides,
};
}
function resolveBoundedNumber(value, fallback, min, max) {
const parsed = Number(value);
if (!Number.isFinite(parsed)) return fallback;
return Math.min(max, Math.max(min, parsed));
}
function normalizeAgentInjectionMode(value) {
const mode = String(value ?? 'off').trim().toLowerCase();
return ['off', 'shadow', 'canary', 'active'].includes(mode) ? mode : 'off';
}
function normalizeUserIdList(value) {
return [...new Set(String(value ?? '')
.split(/[\s,]+/u)
.map((item) => item.trim())
.filter(Boolean))].slice(0, 1000);
}
export function createLegacyMemoryBackend(conversationMemoryService) {
return {
name: 'legacy-conversation-memory',
isAvailable() {
return Boolean(conversationMemoryService);
},
async resolve({ userId, limit = 40 } = {}) {
if (!conversationMemoryService?.listMemories || !userId) {
return { memories: [] };
}
const memories = await conversationMemoryService.listMemories(userId, { limit });
return { memories };
},
async write({ userId, sessionId, messages = [] } = {}) {
if (!conversationMemoryService?.saveAndAnalyze || !userId || !sessionId) {
return { saved: 0, analyzed: 0, memories: 0 };
}
return conversationMemoryService.saveAndAnalyze(sessionId, userId, messages);
},
async compact({ userId } = {}) {
if (!conversationMemoryService?.analyzeUser || !userId) {
return { analyzed: 0, memories: 0 };
}
return conversationMemoryService.analyzeUser(userId);
},
};
}
function normalizeMemoryItem(item) {
if (typeof item === 'string') {
const text = normalizeText(item);
return text ? { label: 'memory', text } : null;
}
if (!item || typeof item !== 'object') return null;
const text = normalizeText(
item.text ?? item.content ?? item.memoryText ?? item.memory_text ?? '',
);
if (!text) return null;
const normalized = {
label: normalizeText(item.label ?? item.type ?? 'memory', 80) ?? 'memory',
text,
};
if (item.id != null) normalized.id = String(item.id);
const hasCreatedAt = Object.prototype.hasOwnProperty.call(item, 'createdAt')
|| Object.prototype.hasOwnProperty.call(item, 'created_at');
if (hasCreatedAt) normalized.createdAt = item.createdAt ?? item.created_at ?? null;
const score = item.score ?? item.confidence;
if (score != null && Number.isFinite(Number(score))) normalized.score = Number(score);
return normalized;
}
function normalizeStringList(values) {
if (!Array.isArray(values)) return [];
return values
.map((value) => (
typeof value === 'string'
? normalizeText(value)
: normalizeText(value?.text ?? value?.content ?? value?.summary ?? '')
))
.filter(Boolean)
.slice(0, MAX_CONTEXT_ITEMS);
}
function normalizeMemoryList(values) {
if (!Array.isArray(values)) return [];
return values
.map((item) => normalizeMemoryItem(item))
.filter(Boolean)
.slice(0, MAX_CONTEXT_ITEMS);
}
function normalizeResolvePayload(payload, source, policy) {
const memories = normalizeMemoryList(payload?.memories);
const semanticMemories = Array.isArray(payload?.semanticMemories)
? normalizeStringList(payload.semanticMemories)
: memories;
const contextGoals = normalizeStringList(
Array.isArray(payload?.contextGoals)
? payload.contextGoals
: (Array.isArray(payload?.userLongTermGoals) ? payload.userLongTermGoals : payload?.activeGoals),
);
return {
ok: true,
enabled: true,
skipped: false,
reason: null,
degraded: false,
source,
profile: policy.profileEnabled ? (payload?.profile ?? null) : null,
semanticMemories,
behaviorSummary: normalizeText(payload?.behaviorSummary, 4000),
activeGoals: contextGoals,
contextGoals,
memories,
};
}
function supportsOperation(backend, operation) {
if (!operation) return true;
return typeof backend?.[operation] === 'function';
}
function isBackendAvailable(backend) {
try {
return backend?.isAvailable?.() !== false;
} catch {
return false;
}
}
function selectBackend(backends, policy, operation = null) {
if (!Array.isArray(backends) || backends.length === 0) return null;
const preferred = policy.backend;
if (preferred && preferred !== 'legacy') {
const exact = backends.find((backend) => backend?.name === preferred);
if (exact && isBackendAvailable(exact) && supportsOperation(exact, operation)) {
return exact;
}
}
return backends.find((backend) => (
isBackendAvailable(backend) && supportsOperation(backend, operation)
)) ?? null;
}
function backendStatus(backend) {
const name = String(backend?.name ?? 'unknown');
let available = true;
try {
available = backend?.isAvailable?.() !== false;
} catch {
available = false;
}
const status = {
name,
available,
supports: {
resolve: typeof backend?.resolve === 'function',
write: typeof backend?.write === 'function',
compact: typeof backend?.compact === 'function',
},
};
if (backend?.category) status.category = String(backend.category);
if (backend?.role) status.role = String(backend.role);
if (backend?.flag) status.flag = String(backend.flag);
const reason = typeof backend?.getUnavailableReason === 'function'
? backend.getUnavailableReason()
: backend?.unavailableReason;
if (!available && reason) status.reason = String(reason);
return status;
}
export function createMemoryV2({
legacyMemoryService = null,
backends = null,
policy = null,
env = process.env,
logger = console,
personalShadowPipeline = null,
} = {}) {
const resolvedPolicy = policy ?? resolveMemoryV2Policy({ env });
const resolvedBackends = backends ?? [
createLegacyMemoryBackend(legacyMemoryService),
...createMemoryV2PluginBackends(),
];
const shadowPipeline = personalShadowPipeline ?? createPersonalMemoryShadowPipeline({ env });
async function failOpen(operation, err, fallback) {
logger?.warn?.(
`[memory-v2] ${operation} degraded: ${err instanceof Error ? err.message : err}`,
);
if (resolvedPolicy.failOpen) return fallback;
throw err;
}
async function resolve(input = {}) {
if (!resolvedPolicy.enabled) {
return emptyResolveResult({ enabled: false, skipped: true, reason: 'disabled' });
}
const backend = selectBackend(resolvedBackends, resolvedPolicy, 'resolve');
if (!backend?.resolve) {
return emptyResolveResult({ enabled: true, skipped: true, reason: 'no_backend' });
}
try {
const payload = await backend.resolve(input);
return normalizeResolvePayload(payload, backend.name ?? 'unknown', resolvedPolicy);
} catch (err) {
return failOpen(
'resolve',
err,
emptyResolveResult({
enabled: true,
degraded: true,
reason: 'backend_failure',
source: backend.name ?? 'unknown',
}),
);
}
}
async function write(input = {}) {
if (!resolvedPolicy.enabled) return skippedWriteResult('disabled');
if (!resolvedPolicy.eventLogEnabled) return skippedWriteResult('event_log_disabled');
const backend = selectBackend(resolvedBackends, resolvedPolicy, 'write');
if (!backend?.write) return skippedWriteResult('no_backend');
try {
const result = await backend.write(input);
if (shadowPipeline?.config?.enabled) {
void Promise.resolve(shadowPipeline.observeWrite(input)).catch((err) => {
logger?.warn?.(
`[memory-v2] personal shadow pipeline skipped: ${err instanceof Error ? err.message : err}`,
);
});
}
return {
ok: true,
enabled: true,
skipped: false,
reason: null,
source: backend.name ?? 'unknown',
saved: Number(result?.saved ?? 0),
analyzed: Number(result?.analyzed ?? 0),
memories: Number(result?.memories ?? 0),
};
} catch (err) {
return failOpen('write', err, {
ok: false,
enabled: true,
skipped: true,
reason: 'backend_failure',
source: backend.name ?? 'unknown',
saved: 0,
analyzed: 0,
memories: 0,
});
}
}
async function compact(input = {}) {
if (!resolvedPolicy.enabled) {
return { ok: true, enabled: false, skipped: true, reason: 'disabled' };
}
const backend = selectBackend(resolvedBackends, resolvedPolicy, 'compact');
if (!backend?.compact) {
return { ok: true, enabled: true, skipped: true, reason: 'no_backend' };
}
try {
const result = await backend.compact(input);
return {
ok: true,
enabled: true,
skipped: false,
reason: null,
source: backend.name ?? 'unknown',
analyzed: Number(result?.analyzed ?? 0),
memories: Number(result?.memories ?? 0),
};
} catch (err) {
return failOpen('compact', err, {
ok: false,
enabled: true,
skipped: true,
reason: 'backend_failure',
source: backend.name ?? 'unknown',
analyzed: 0,
memories: 0,
});
}
}
async function observePersonalMemory(input = {}) {
if (!shadowPipeline?.config?.enabled) {
return { enabled: false, skipped: true, reason: 'disabled' };
}
try {
return await shadowPipeline.observeWrite(input);
} catch (err) {
logger?.warn?.(
`[memory-v2] personal shadow observation skipped: ${err instanceof Error ? err.message : err}`,
);
return {
enabled: true,
skipped: true,
reason: 'pipeline_failure',
};
}
}
function getStatus() {
const backends = resolvedBackends.map((backend) => backendStatus(backend));
const selected = selectBackend(resolvedBackends, resolvedPolicy, 'resolve');
const status = {
enabled: Boolean(resolvedPolicy.enabled),
backend: resolvedPolicy.backend,
selectedBackend: selected?.name ?? null,
profileEnabled: Boolean(resolvedPolicy.profileEnabled),
eventLogEnabled: Boolean(resolvedPolicy.eventLogEnabled),
vectorEnabled: Boolean(resolvedPolicy.vectorEnabled),
failOpen: Boolean(resolvedPolicy.failOpen),
runtimeControl: {
agentResolveEnabled: Boolean(resolvedPolicy.agentResolveEnabled),
agentInjectionMode: resolvedPolicy.agentInjectionMode,
agentCanaryUserIds: resolvedPolicy.agentCanaryUserIds,
agentResolveLimit: Number(resolvedPolicy.agentResolveLimit),
agentResolveTimeoutMs: Number(resolvedPolicy.agentResolveTimeoutMs),
promotionEnabled: Boolean(resolvedPolicy.promotionEnabled),
compactionV2Enabled: Boolean(resolvedPolicy.compactionV2Enabled),
reflectionEnabled: Boolean(resolvedPolicy.reflectionEnabled),
lifecycleWorkerEnabled: Boolean(resolvedPolicy.lifecycleWorkerEnabled),
lifecycleRolloutMode: resolvedPolicy.lifecycleRolloutMode,
lifecycleRolloutUserIds: resolvedPolicy.lifecycleRolloutUserIds,
},
backends,
};
if (shadowPipeline?.config?.enabled) {
status.personalMemory = shadowPipeline.getStatus();
}
return status;
}
return {
policy: resolvedPolicy,
getStatus,
resolve,
write,
compact,
observePersonalMemory,
};
}