openclaw/extensions/active-memory/transcript-watch.ts
Peter Steinberger a72ff7a93f
fix(memory): preserve rotated recall and simplify consolidation (#137967)
* refactor(memory): simplify recall and consolidation owners

* test(memory): isolate recall state and timeout clocks

* test(memory): await callback transcript before timing out
2026-09-04 02:18:00 -07:00

183 lines
5.6 KiB
TypeScript

import { normalizeLowercaseStringOrEmpty } from "openclaw/plugin-sdk/string-coerce-runtime";
import {
extractActiveMemorySearchDebugFromSessionRecord,
extractToolResultNameFromSessionRecord,
hasTerminalUnavailableMemoryResultInSessionRecord,
hasUnavailableMemoryResultInSessionRecord,
hasUsableMemoryResultInSessionRecord,
streamActiveMemoryTranscriptRecords,
} from "./transcript.js";
import {
TERMINAL_MEMORY_SEARCH_POLL_INTERVAL_MS,
type ActiveMemorySearchDebug,
type ActiveMemoryTranscriptSource,
type TerminalMemorySearchResult,
type TerminalMemorySearchWatch,
type TranscriptReadLimits,
} from "./types.js";
async function readMergedActiveMemoryTranscriptState(params: {
sources: readonly ActiveMemoryTranscriptSource[];
toolsAllow: readonly string[];
}): Promise<{
searchDebug?: ActiveMemorySearchDebug;
hasUsableMemoryResult: boolean;
hasUnavailableMemorySearchResult: boolean;
}> {
let searchDebug: ActiveMemorySearchDebug | undefined;
let hasUsableMemoryResult = false;
let hasUnavailableMemorySearchResult = false;
for (const source of params.sources) {
await streamActiveMemoryTranscriptRecords({
source,
onRecord: (record) => {
searchDebug = extractActiveMemorySearchDebugFromSessionRecord(record) ?? searchDebug;
hasUnavailableMemorySearchResult ||= hasUnavailableMemoryResultInSessionRecord(
record,
params.toolsAllow,
);
hasUsableMemoryResult ||= hasUsableMemoryResultInSessionRecord(record, params.toolsAllow);
},
});
}
return { searchDebug, hasUsableMemoryResult, hasUnavailableMemorySearchResult };
}
async function readTerminalMemorySearchResult(
source: ActiveMemoryTranscriptSource,
limits?: TranscriptReadLimits,
toolsAllow?: readonly string[],
): Promise<TerminalMemorySearchResult | undefined> {
// memory_get consumes a path discovered by another tool; it is not an
// independent fallback that should delay terminal unavailability.
const recallPathNames = new Set(
toolsAllow
?.map((toolName) => normalizeLowercaseStringOrEmpty(toolName))
.filter((toolName) => toolName && toolName !== "memory_get"),
);
if (recallPathNames.size === 0) {
return undefined;
}
const unavailablePathNames = new Set<string>();
let hasUsableMemoryResult = false;
let searchDebug: ActiveMemorySearchDebug | undefined;
await streamActiveMemoryTranscriptRecords({
source,
limits,
onRecord: (record) => {
hasUsableMemoryResult ||= hasUsableMemoryResultInSessionRecord(record, toolsAllow);
searchDebug = extractActiveMemorySearchDebugFromSessionRecord(record) ?? searchDebug;
const toolName = extractToolResultNameFromSessionRecord(record);
if (!toolName || !recallPathNames.has(toolName)) {
return false;
}
if (hasTerminalUnavailableMemoryResultInSessionRecord(record, toolsAllow ?? [])) {
unavailablePathNames.add(toolName);
} else {
unavailablePathNames.delete(toolName);
}
return false;
},
});
if (unavailablePathNames.size !== recallPathNames.size) {
return undefined;
}
return {
status: "unavailable",
hasUsableMemoryResult,
searchDebug,
};
}
async function readTerminalMemorySearchResultFromSources(
sources: readonly ActiveMemoryTranscriptSource[],
limits: TranscriptReadLimits | undefined,
toolsAllow: readonly string[],
): Promise<TerminalMemorySearchResult | undefined> {
for (const source of sources) {
const result = await readTerminalMemorySearchResult(source, limits, toolsAllow);
if (result) {
return result;
}
}
return undefined;
}
function watchTerminalMemorySearchResult(params: {
getTranscriptSources: () => readonly ActiveMemoryTranscriptSource[];
abortSignal: AbortSignal;
toolsAllow: readonly string[];
}): TerminalMemorySearchWatch {
let stopped = false;
let timeoutId: ReturnType<typeof setTimeout> | undefined;
let inFlight = false;
let resolveWatch: (result: TerminalMemorySearchResult) => void = () => {};
const stop = () => {
if (stopped) {
return;
}
stopped = true;
if (timeoutId) {
clearTimeout(timeoutId);
timeoutId = undefined;
}
params.abortSignal.removeEventListener("abort", onAbort);
};
const finish = (result: TerminalMemorySearchResult) => {
stop();
resolveWatch(result);
};
const schedule = () => {
if (stopped) {
return;
}
timeoutId = setTimeout(() => {
void tick();
}, TERMINAL_MEMORY_SEARCH_POLL_INTERVAL_MS);
timeoutId.unref?.();
};
const tick = async () => {
if (stopped || inFlight) {
return;
}
if (params.abortSignal.aborted) {
stop();
return;
}
inFlight = true;
try {
const result = await readTerminalMemorySearchResultFromSources(
params.getTranscriptSources(),
undefined,
params.toolsAllow,
);
// Execution can settle while this transcript read is still in flight.
if (stopped || params.abortSignal.aborted) {
return;
}
if (result) {
finish(result);
return;
}
} catch {
// Transcript polling is opportunistic; normal timeout handling remains authoritative.
} finally {
inFlight = false;
}
schedule();
};
function onAbort() {
stop();
}
const promise = new Promise<TerminalMemorySearchResult>((resolve) => {
resolveWatch = resolve;
params.abortSignal.addEventListener("abort", onAbort, { once: true });
void tick();
});
return {
promise,
stop,
};
}
export { readMergedActiveMemoryTranscriptState, watchTerminalMemorySearchResult };