Skip to content

Commit 0c90ccb

Browse files
committed
refactor: add memory and QMD session identity mapping
1 parent 97b97a9 commit 0c90ccb

22 files changed

Lines changed: 1403 additions & 129 deletions

extensions/memory-core/src/memory/manager-session-reindex.ts

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,19 +1,20 @@
11
// Memory Core plugin module implements manager session reindex behavior.
2+
import type { MemorySyncParams } from "openclaw/plugin-sdk/memory-core-host-engine-storage";
3+
24
export function shouldSyncSessionsForReindex(params: {
35
hasSessionSource: boolean;
46
sessionsDirty: boolean;
57
sessionsFullRetryDirty?: boolean;
68
dirtySessionFileCount: number;
7-
sync?: {
8-
reason?: string;
9-
force?: boolean;
10-
sessionFiles?: string[];
11-
};
9+
sync?: MemorySyncParams;
1210
needsFullReindex?: boolean;
1311
}): boolean {
1412
if (!params.hasSessionSource) {
1513
return false;
1614
}
15+
if (params.sync?.sessions?.some((session) => session.sessionId.trim().length > 0)) {
16+
return true;
17+
}
1718
if (params.sync?.sessionFiles?.some((sessionFile) => sessionFile.trim().length > 0)) {
1819
return true;
1920
}

extensions/memory-core/src/memory/manager-session-sync-state.test.ts

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -66,6 +66,19 @@ describe("memory session sync state", () => {
6666
expect(plan.activePaths).toEqual(new Set(["sessions/incremental.jsonl"]));
6767
});
6868

69+
it("marks identity-targeted syncs as session work", async () => {
70+
const { shouldSyncSessionsForReindex } = await import("./manager-session-reindex.js");
71+
72+
expect(
73+
shouldSyncSessionsForReindex({
74+
hasSessionSource: true,
75+
sessionsDirty: false,
76+
dirtySessionFileCount: 0,
77+
sync: { sessions: [{ agentId: "main", sessionId: "targeted" }] },
78+
}),
79+
).toBe(true);
80+
});
81+
6982
it("marks missing and changed startup session files dirty", () => {
7083
const dirtyFiles = resolveMemorySessionStartupDirtyFiles({
7184
files: [

extensions/memory-core/src/memory/manager-sync-control.ts

Lines changed: 47 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,11 @@
22
import type { DatabaseSync } from "node:sqlite";
33
import { formatErrorMessage } from "openclaw/plugin-sdk/error-runtime";
44
import { createSubsystemLogger } from "openclaw/plugin-sdk/memory-core-host-engine-foundation";
5-
import type { MemorySyncProgressUpdate } from "openclaw/plugin-sdk/memory-core-host-engine-storage";
5+
import type {
6+
MemorySessionSyncTarget,
7+
MemorySyncParams,
8+
MemorySyncProgressUpdate,
9+
} from "openclaw/plugin-sdk/memory-core-host-engine-storage";
610

711
const log = createSubsystemLogger("memory");
812

@@ -19,6 +23,7 @@ export type MemoryReadonlyRecoveryState = {
1923
runSync: (params?: {
2024
reason?: string;
2125
force?: boolean;
26+
sessions?: MemorySessionSyncTarget[];
2227
sessionFiles?: string[];
2328
progress?: (update: MemorySyncProgressUpdate) => void;
2429
}) => Promise<void>;
@@ -80,12 +85,7 @@ export function extractMemoryErrorReason(err: unknown): string {
8085

8186
export async function runMemorySyncWithReadonlyRecovery(
8287
state: MemoryReadonlyRecoveryState,
83-
params?: {
84-
reason?: string;
85-
force?: boolean;
86-
sessionFiles?: string[];
87-
progress?: (update: MemorySyncProgressUpdate) => void;
88-
},
88+
params?: MemorySyncParams,
8989
): Promise<void> {
9090
try {
9191
await state.runSync(params);
@@ -121,37 +121,46 @@ export function enqueueMemoryTargetedSessionSync(
121121
isClosed: () => boolean;
122122
getSyncing: () => Promise<void> | null;
123123
getQueuedSessionFiles: () => Set<string>;
124+
getQueuedSessions: () => Map<string, MemorySessionSyncTarget>;
124125
getQueuedSessionSync: () => Promise<void> | null;
125126
setQueuedSessionSync: (value: Promise<void> | null) => void;
126-
sync: (params?: {
127-
reason?: string;
128-
force?: boolean;
129-
sessionFiles?: string[];
130-
progress?: (update: MemorySyncProgressUpdate) => void;
131-
}) => Promise<void>;
127+
sync: (params?: MemorySyncParams) => Promise<void>;
132128
},
133-
sessionFiles?: string[],
129+
targets?: Pick<MemorySyncParams, "sessions" | "sessionFiles">,
134130
): Promise<void> {
135131
const queuedSessionFiles = state.getQueuedSessionFiles();
136-
for (const sessionFile of sessionFiles ?? []) {
132+
for (const sessionFile of targets?.sessionFiles ?? []) {
137133
const trimmed = sessionFile.trim();
138134
if (trimmed) {
139135
queuedSessionFiles.add(trimmed);
140136
}
141137
}
142-
if (queuedSessionFiles.size === 0) {
138+
const queuedSessions = state.getQueuedSessions();
139+
for (const session of targets?.sessions ?? []) {
140+
const normalized = normalizeQueuedMemorySessionSyncTarget(session);
141+
if (normalized) {
142+
queuedSessions.set(memorySessionSyncTargetKey(normalized), normalized);
143+
}
144+
}
145+
if (queuedSessionFiles.size === 0 && queuedSessions.size === 0) {
143146
return state.getSyncing() ?? Promise.resolve();
144147
}
145148
if (!state.getQueuedSessionSync()) {
146149
state.setQueuedSessionSync(
147150
(async () => {
148151
try {
149152
await state.getSyncing()?.catch(() => undefined);
150-
while (!state.isClosed() && state.getQueuedSessionFiles().size > 0) {
153+
while (
154+
!state.isClosed() &&
155+
(state.getQueuedSessionFiles().size > 0 || state.getQueuedSessions().size > 0)
156+
) {
151157
const pendingSessionFiles = Array.from(state.getQueuedSessionFiles());
158+
const pendingSessions = Array.from(state.getQueuedSessions().values());
152159
state.getQueuedSessionFiles().clear();
160+
state.getQueuedSessions().clear();
153161
await state.sync({
154-
reason: "queued-session-files",
162+
reason: "queued-sessions",
163+
sessions: pendingSessions,
155164
sessionFiles: pendingSessionFiles,
156165
});
157166
}
@@ -163,3 +172,23 @@ export function enqueueMemoryTargetedSessionSync(
163172
}
164173
return state.getQueuedSessionSync() ?? Promise.resolve();
165174
}
175+
176+
function normalizeQueuedMemorySessionSyncTarget(
177+
target: MemorySessionSyncTarget,
178+
): MemorySessionSyncTarget | null {
179+
const sessionId = target.sessionId.trim();
180+
if (!sessionId) {
181+
return null;
182+
}
183+
const agentId = target.agentId?.trim();
184+
const sessionKey = target.sessionKey?.trim();
185+
return {
186+
...(agentId ? { agentId } : {}),
187+
sessionId,
188+
...(sessionKey ? { sessionKey } : {}),
189+
};
190+
}
191+
192+
function memorySessionSyncTargetKey(target: MemorySessionSyncTarget): string {
193+
return [target.agentId ?? "", target.sessionId, target.sessionKey ?? ""].join("\0");
194+
}

extensions/memory-core/src/memory/manager-sync-ops.startup-catchup.test.ts

Lines changed: 56 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import {
1010
} from "openclaw/plugin-sdk/memory-core-host-engine-foundation";
1111
import type {
1212
MemorySource,
13+
MemorySyncParams,
1314
MemorySyncProgressUpdate,
1415
} from "openclaw/plugin-sdk/memory-core-host-engine-storage";
1516
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
@@ -27,6 +28,7 @@ type MemoryIndexEntry = {
2728
type SyncParams = {
2829
reason?: string;
2930
force?: boolean;
31+
sessions?: MemorySyncParams["sessions"];
3032
sessionFiles?: string[];
3133
progress?: (update: MemorySyncProgressUpdate) => void;
3234
};
@@ -38,14 +40,33 @@ class SessionStartupCatchupHarness extends MemoryManagerSyncOps {
3840
protected readonly agentId = "main";
3941
protected readonly workspaceDir = "/tmp/openclaw-test-workspace";
4042
protected readonly settings = {
43+
chunking: {
44+
overlap: 0,
45+
tokens: 256,
46+
},
47+
extraPaths: [],
48+
multimodal: {
49+
enabled: false,
50+
modalities: [],
51+
maxFileBytes: 0,
52+
},
53+
provider: "none",
54+
store: {
55+
fts: {
56+
tokenizer: "unicode61",
57+
},
58+
vector: {
59+
enabled: false,
60+
},
61+
},
4162
sync: {
4263
sessions: {
4364
deltaBytes: 100_000,
4465
deltaMessages: 50,
4566
postCompactionForce: true,
4667
},
4768
},
48-
} as ResolvedMemorySearchConfig;
69+
} as unknown as ResolvedMemorySearchConfig;
4970
protected readonly batch = {
5071
enabled: false,
5172
wait: false,
@@ -60,6 +81,7 @@ class SessionStartupCatchupHarness extends MemoryManagerSyncOps {
6081
protected db: DatabaseSync;
6182

6283
readonly syncCalls: SyncParams[] = [];
84+
readonly indexedPaths: string[] = [];
6385

6486
constructor(sourceRows: SourceStateRow[]) {
6587
super();
@@ -81,6 +103,10 @@ class SessionStartupCatchupHarness extends MemoryManagerSyncOps {
81103
return await this.markSessionStartupCatchupDirtyFiles();
82104
}
83105

106+
async runSyncForTest(params?: MemorySyncParams): Promise<void> {
107+
await this.runSync(params);
108+
}
109+
84110
getDirtySessionFiles(): string[] {
85111
return Array.from(this.sessionsDirtyFiles);
86112
}
@@ -97,7 +123,7 @@ class SessionStartupCatchupHarness extends MemoryManagerSyncOps {
97123
return [];
98124
}
99125

100-
protected async sync(params?: SyncParams): Promise<void> {
126+
protected async sync(params?: MemorySyncParams): Promise<void> {
101127
this.syncCalls.push(params ?? {});
102128
}
103129

@@ -120,9 +146,11 @@ class SessionStartupCatchupHarness extends MemoryManagerSyncOps {
120146
protected assertRequiredProviderAvailable(): void {}
121147

122148
protected async indexFile(
123-
_entry: MemoryIndexEntry,
149+
entry: MemoryIndexEntry,
124150
_options: { source: MemorySource; content?: string },
125-
): Promise<void> {}
151+
): Promise<void> {
152+
this.indexedPaths.push(entry.path);
153+
}
126154
}
127155

128156
describe("session startup catch-up", () => {
@@ -304,4 +332,28 @@ describe("session startup catch-up", () => {
304332
openSpy.mockRestore();
305333
}
306334
});
335+
336+
it("does not fall back to full session sync when identity targets normalize away", async () => {
337+
await writeSessionFile("thread.jsonl");
338+
const harness = new SessionStartupCatchupHarness([]);
339+
340+
await harness.runSyncForTest({
341+
reason: "queued-sessions",
342+
sessions: [{ agentId: "other", sessionId: "thread" }],
343+
});
344+
345+
expect(harness.indexedPaths).toEqual([]);
346+
});
347+
348+
it("does not fall back to full session sync for malformed identity session ids", async () => {
349+
await writeSessionFile("thread.jsonl");
350+
const harness = new SessionStartupCatchupHarness([]);
351+
352+
await harness.runSyncForTest({
353+
reason: "queued-sessions",
354+
sessions: [{ agentId: "main", sessionId: "bad/nested" }],
355+
});
356+
357+
expect(harness.indexedPaths).toEqual([]);
358+
});
307359
});

0 commit comments

Comments
 (0)