Skip to content
Draft
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
79 changes: 62 additions & 17 deletions dist/index.js
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ import { buildReflectionMappedMetadata } from "./src/reflection-mapped-metadata.
import { gateMappedReflectionEntries } from "./src/reflection-mapped-admission.js";
import { createMemoryCLI } from "./cli.js";
import { isNoise } from "./src/noise-filter.js";
import { normalizeAutoCaptureText } from "./src/auto-capture-cleanup.js";
import { buildConversationTurnsForExtraction, dedupePairWindow, formatConversationTranscript, normalizeAutoCaptureText, trimTranscriptToTagBoundary, trimTurnsToUserCap, } from "./src/auto-capture-cleanup.js";
// Import smart extraction & lifecycle components
import { SmartExtractor, createExtractionRateLimiter } from "./src/smart-extractor.js";
import { compressTexts, estimateConversationValue } from "./src/session-compressor.js";
Expand Down Expand Up @@ -945,7 +945,7 @@ function extractTextFromToolResult(result) {
return "";
}
}
function summarizeRecentConversationMessages(messages, messageCount) {
function summarizeRecentConversationMessages(messages, messageCount, format = "tagged") {
if (!Array.isArray(messages) || messages.length === 0)
return null;
const recent = [];
Expand All @@ -960,14 +960,17 @@ function summarizeRecentConversationMessages(messages, messageCount) {
const text = extractTextContent(msg.content);
if (!text || shouldSkipReflectionMessage(role, text))
continue;
recent.push(`${role}: ${redactSecrets(text)}`);
recent.push({ role, text: redactSecrets(text) });
}
if (recent.length === 0)
return null;
recent.reverse();
return recent.join("\n");
if (format === "labeled") {
return recent.map((turn) => `${turn.role}: ${turn.text}`).join("\n");
}
return formatConversationTranscript(recent);
}
async function readSessionConversationForReflection(filePath, messageCount) {
async function readSessionConversationForReflection(filePath, messageCount, format = "tagged") {
try {
const lines = (await readFile(filePath, "utf-8")).trim().split("\n");
const messages = [];
Expand All @@ -982,14 +985,14 @@ async function readSessionConversationForReflection(filePath, messageCount) {
// ignore JSON parse errors
}
}
return summarizeRecentConversationMessages(messages, messageCount);
return summarizeRecentConversationMessages(messages, messageCount, format);
}
catch {
return null;
}
}
export async function readSessionConversationWithResetFallback(sessionFilePath, messageCount) {
const primary = await readSessionConversationForReflection(sessionFilePath, messageCount);
export async function readSessionConversationWithResetFallback(sessionFilePath, messageCount, format = "tagged") {
const primary = await readSessionConversationForReflection(sessionFilePath, messageCount, format);
if (primary)
return primary;
try {
Expand All @@ -999,7 +1002,7 @@ export async function readSessionConversationWithResetFallback(sessionFilePath,
const resetCandidates = await sortFileNamesByMtimeDesc(dir, files.filter((name) => name.startsWith(resetPrefix)));
if (resetCandidates.length > 0) {
const latestResetPath = join(dir, resetCandidates[0]);
return await readSessionConversationForReflection(latestResetPath, messageCount);
return await readSessionConversationForReflection(latestResetPath, messageCount, format);
}
}
catch {
Expand All @@ -1016,7 +1019,7 @@ async function ensureDailyLogFile(dailyPath, dateStr) {
}
}
export function buildReflectionPrompt(conversation, maxInputChars, toolErrorSignals = []) {
const clipped = conversation.slice(-maxInputChars);
const clipped = trimTranscriptToTagBoundary(conversation, maxInputChars);
const errorHints = toolErrorSignals.length > 0
? toolErrorSignals
.map((e, i) => `${i + 1}. [${e.toolName}] ${e.summary} (sig:${e.signatureHash.slice(0, 8)})`)
Expand All @@ -1025,6 +1028,10 @@ export function buildReflectionPrompt(conversation, maxInputChars, toolErrorSign
return [
"You are generating a durable MEMORY REFLECTION entry for an AI assistant system.",
"",
"The INPUT transcript is a sequence of tagged blocks in chronological order:",
"- <user_message>...</user_message> wraps ONE message written by the human user.",
"- <assistant_message>...</assistant_message> wraps ONE message written by the AI assistant.",
"",
"Output Markdown only. No intro text. No outro text. No extra headings.",
"",
"Use these headings exactly once, in this exact order, with exact spelling:",
Expand Down Expand Up @@ -1124,9 +1131,7 @@ export function buildReflectionPrompt(conversation, maxInputChars, toolErrorSign
errorHints,
"",
"INPUT:",
"```",
clipped,
"```",
].join("\n");
}
function buildReflectionFallbackText() {
Expand Down Expand Up @@ -1848,6 +1853,7 @@ function _initPluginState(api) {
const admissionRejectionAuditWriter = createAdmissionRejectionAuditWriter(config, resolvedDbPath, api);
smartExtractor = new SmartExtractor(store, embedder, llmClient, {
user: "User",
captureAssistantEligible: config.captureAssistant === true,
extractMinMessages: config.extractMinMessages ?? 4,
extractMaxChars: config.extractMaxChars ?? 8000,
defaultScope: config.scopes?.default ?? "global",
Expand Down Expand Up @@ -1887,6 +1893,7 @@ function _initPluginState(api) {
const autoCaptureSeenTextCount = new Map();
const autoCapturePendingIngressTexts = new Map();
const autoCaptureRecentTexts = new Map();
const autoCaptureRecentPairTurns = new Map();
return {
config,
resolvedDbPath,
Expand Down Expand Up @@ -1914,6 +1921,7 @@ function _initPluginState(api) {
autoCaptureSeenTextCount,
autoCapturePendingIngressTexts,
autoCaptureRecentTexts,
autoCaptureRecentPairTurns,
};
}
export function isAgentOrSessionExcluded(agentId, sessionKey, patterns) {
Expand Down Expand Up @@ -2022,7 +2030,7 @@ const memoryLanceDBProPlugin = {
_registeredApisMap.delete(api); // dual-track rollback: Map un-claim
throw err;
}
const { config, resolvedDbPath, vectorDim, store, embedder, retriever, canonicalCorpusIndexer, dreamingEngine, dreamingScheduler, scopeManager, migrator, smartExtractor, mdMirror, decayEngine, tierManager, extractionRateLimiter, reflectionErrorStateBySession, reflectionDerivedBySession, reflectionDerivedSuppressionBySession, reflectionByAgentCache, reflectionByAgentCacheGeneration, recallHistory, turnCounter, autoCaptureSeenTextCount, autoCapturePendingIngressTexts, autoCaptureRecentTexts, } = singleton;
const { config, resolvedDbPath, vectorDim, store, embedder, retriever, canonicalCorpusIndexer, dreamingEngine, dreamingScheduler, scopeManager, migrator, smartExtractor, mdMirror, decayEngine, tierManager, extractionRateLimiter, reflectionErrorStateBySession, reflectionDerivedBySession, reflectionDerivedSuppressionBySession, reflectionByAgentCache, reflectionByAgentCacheGeneration, recallHistory, turnCounter, autoCaptureSeenTextCount, autoCapturePendingIngressTexts, autoCaptureRecentTexts, autoCaptureRecentPairTurns, } = singleton;
warnForDisabledChannelPlugin(api.config, api.logger);
async function sleep(ms, signal) {
if (signal?.aborted) {
Expand Down Expand Up @@ -2903,8 +2911,10 @@ const memoryLanceDBProPlugin = {
: scopeManager.getDefaultScope(agentId);
const sessionKey = ctx?.sessionKey || event.sessionKey || "unknown";
api.logger.debug(`memory-lancedb-pro: auto-capture agent_end payload for agent ${agentId} (sessionKey=${sessionKey}, captureAssistant=${config.captureAssistant === true}, ${summarizeAgentEndMessages(event.messages)})`);
// Extract text content from messages
// Extract text content from messages, keeping the role-tagged
// message-loop order alongside the flat eligible-text list.
const eligibleTexts = [];
const messageLoopTurns = [];
let skippedAutoCaptureTexts = 0;
for (const msg of event.messages) {
if (!msg || typeof msg !== "object") {
Expand All @@ -2925,6 +2935,7 @@ const memoryLanceDBProPlugin = {
}
else {
eligibleTexts.push(normalized);
messageLoopTurns.push({ role: role, text: normalized });
}
continue;
}
Expand All @@ -2943,6 +2954,7 @@ const memoryLanceDBProPlugin = {
}
else {
eligibleTexts.push(normalized);
messageLoopTurns.push({ role: role, text: normalized });
}
}
}
Expand Down Expand Up @@ -2979,6 +2991,29 @@ const memoryLanceDBProPlugin = {
autoCaptureRecentTexts.set(sessionKey, nextRecentTexts);
pruneMapIfOver(autoCaptureRecentTexts, AUTO_CAPTURE_MAP_MAX_ENTRIES);
}
// Rolling PAIR window sized by autoCaptureContextTurns (0 =
// disabled: each extraction sees only its own call's turns, and
// nothing is retained between calls). When enabled, this call's
// new pairs -- kept user turns with the assistant replies
// interleaved in true order -- extend what earlier calls
// buffered, bounded to autoCaptureContextTurns user turns (or
// this call's own new-user count when larger, so unextracted
// user turns are never trimmed out of their own transcript).
const thisCallTurns = buildConversationTurnsForExtraction({
messageLoopTurns,
eligibleTexts,
newUserTexts: newTexts,
});
const contextTurns = config.autoCaptureContextTurns ?? 0;
const priorPairTurns = contextTurns > 0 ? autoCaptureRecentPairTurns.get(sessionKey) || [] : [];
const pairWindowTurns = trimTurnsToUserCap(dedupePairWindow([...priorPairTurns, ...thisCallTurns]), Math.max(contextTurns, thisCallTurns.filter((turn) => turn.role === "user").length));
if (contextTurns === 0) {
autoCaptureRecentPairTurns.delete(sessionKey);
}
else if (thisCallTurns.length > 0) {
autoCaptureRecentPairTurns.set(sessionKey, pairWindowTurns);
pruneMapIfOver(autoCaptureRecentPairTurns, AUTO_CAPTURE_MAP_MAX_ENTRIES);
}
const minMessages = config.extractMinMessages ?? 4;
if (skippedAutoCaptureTexts > 0) {
api.logger.debug(`memory-lancedb-pro: auto-capture skipped ${skippedAutoCaptureTexts} injected/system text block(s) for agent ${agentId}`);
Expand Down Expand Up @@ -3035,10 +3070,14 @@ const memoryLanceDBProPlugin = {
if (cumulativeCount >= minMessages) {
api.logger.debug(`memory-lancedb-pro: auto-capture running smart extraction for agent ${agentId} (cumulative=${cumulativeCount} >= minMessages=${minMessages}, cleanTexts=${cleanTexts.length})`);
const conversationText = cleanTexts.join("\n");
// The pair window is the transcript; user turns the noise
// filter dropped stay out of it so they cannot become sources.
const noiseDroppedTexts = new Set(texts.filter((text) => !cleanTexts.includes(text)));
const finalConversationTurns = pairWindowTurns.filter((turn) => !(turn.role === "user" && noiseDroppedTexts.has(turn.text)));
// issue #417 Fix #10: prevent hook crash on LLM API errors / network timeouts
let stats = null;
try {
stats = await smartExtractor.extractAndPersist(conversationText, sessionKey, { scope: defaultScope, scopeFilter: accessibleScopes, agentId });
stats = await smartExtractor.extractAndPersist(conversationText, sessionKey, { scope: defaultScope, scopeFilter: accessibleScopes, agentId, conversationTurns: finalConversationTurns });
}
catch (err) {
api.logger.error(`memory-lancedb-pro: smart-extract failed for agent ${agentId}: ${String(err)}`);
Expand All @@ -3058,6 +3097,11 @@ const memoryLanceDBProPlugin = {
// consumed history length there instead, so the next turn
// only sees the delta.
autoCaptureSeenTextCount.set(sessionKey, pendingIngressTexts.length > 0 ? 0 : eligibleTexts.length);
// The rolling pair window is deliberately retained across
// successful extractions: deleting it here would mean
// steady-state captures (one extraction per turn) always see
// a bare current pair. The set-time trim bounds it; the
// watermark keeps retained turns from re-becoming sources.
return; // Smart extraction handled everything
}
if ((stats.boundarySkipped ?? 0) === 0) {
Expand Down Expand Up @@ -4183,9 +4227,9 @@ const memoryLanceDBProPlugin = {
return;
}
guard.set(guardKey, now);
const sessionContent = summarizeRecentConversationMessages(event.messages ?? [], sessionMessageCount) ??
const sessionContent = summarizeRecentConversationMessages(event.messages ?? [], sessionMessageCount, "labeled") ??
(typeof event.sessionFile === "string"
? await readSessionConversationWithResetFallback(event.sessionFile, sessionMessageCount)
? await readSessionConversationWithResetFallback(event.sessionFile, sessionMessageCount, "labeled")
: null);
if (!sessionContent) {
guard.delete(guardKey);
Expand Down Expand Up @@ -4779,6 +4823,7 @@ export function parsePluginConfig(value) {
})()
: undefined,
extractMinMessages: parsePositiveInt(cfg.extractMinMessages) ?? 4,
autoCaptureContextTurns: Math.min(10, Math.max(0, Math.floor(Number(cfg.autoCaptureContextTurns)) || 0)),
extractMaxChars: parsePositiveInt(cfg.extractMaxChars) ?? 8000,
scopes: typeof cfg.scopes === "object" && cfg.scopes !== null ? cfg.scopes : undefined,
enableManagementTools: cfg.enableManagementTools === true,
Expand Down
Loading
Loading