mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
fix(memory): keep corpus discovery read-only during forget (#162155)
This commit is contained in:
parent
7dc4a94490
commit
646c42b4c6
13 changed files with 399 additions and 107 deletions
|
|
@ -262,6 +262,10 @@ planning and vector inspection use the same retrieval worker and captured store
|
|||
target. Indexed memory text stays with that reader; the host receives only selected
|
||||
chunk identities, source paths, and counts. Preview remains noncreating, and native
|
||||
vector inspection closes its probe and read-only connection before replying.
|
||||
Forget's corpus discovery requests read-only metadata without unused transcript
|
||||
revisions. Durable session summaries use the existing retained history worker,
|
||||
preserving captured store selection and classification without writable bootstrap.
|
||||
Synchronous corpus callers and process-held transcripts keep their existing owners.
|
||||
Session policy metadata reads and cold bootstrap remain separate. Schemas, stored
|
||||
formats, and update behavior are unchanged.
|
||||
|
||||
|
|
|
|||
|
|
@ -6,7 +6,10 @@ import type { OpenClawConfig } from "openclaw/plugin-sdk/memory-core-host-engine
|
|||
import { encodeMemoryEmbedding } from "openclaw/plugin-sdk/memory-core-host-engine-storage";
|
||||
import { deleteSessionEntry } from "openclaw/plugin-sdk/session-store-runtime";
|
||||
import { appendSessionTranscriptMessageByIdentity } from "openclaw/plugin-sdk/session-transcript-runtime";
|
||||
import { openOpenClawAgentDatabase } from "openclaw/plugin-sdk/sqlite-runtime";
|
||||
import {
|
||||
openOpenClawAgentDatabase,
|
||||
resolveOpenClawAgentSqlitePath,
|
||||
} from "openclaw/plugin-sdk/sqlite-runtime";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import {
|
||||
DREAMING_MEMORY_BACKUP_NAMESPACE,
|
||||
|
|
@ -56,6 +59,48 @@ describe("memory forget", () => {
|
|||
await fixture.cleanup();
|
||||
});
|
||||
|
||||
it("previews an unresolved session without creating the absent agent store", async () => {
|
||||
const databasePath = resolveOpenClawAgentSqlitePath({ agentId: "main" });
|
||||
await expect(fs.stat(databasePath)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
const report = await forgetMemoryEntries({
|
||||
cfg,
|
||||
agentId: "main",
|
||||
sessionIds: ["missing"],
|
||||
dryRun: true,
|
||||
});
|
||||
expect(report).toEqual({
|
||||
agentId: "main",
|
||||
dryRun: true,
|
||||
sessionIds: ["missing"],
|
||||
participantMatches: [],
|
||||
sessionResolutions: [{ sessionId: "missing", source: "unresolved" }],
|
||||
entryKeys: [],
|
||||
mixedLineageEntryKeys: [],
|
||||
untargetableEntryKeys: [],
|
||||
curatedWrites: [],
|
||||
artifacts: {
|
||||
memoryFiles: 0,
|
||||
memoryEntries: 0,
|
||||
memoryLines: 0,
|
||||
sessionCorpusFiles: 0,
|
||||
sessionCorpusLines: 0,
|
||||
indexChunks: 0,
|
||||
indexSources: 0,
|
||||
ftsRows: 0,
|
||||
vectorRows: 0,
|
||||
embeddingCacheRows: 0,
|
||||
shortTermEntries: 0,
|
||||
seenHashScopes: 0,
|
||||
backups: 0,
|
||||
originRows: 0,
|
||||
},
|
||||
refusals: [],
|
||||
});
|
||||
await expect(fs.stat(databasePath)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
await expect(fs.stat(`${databasePath}-wal`)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
await expect(fs.stat(`${databasePath}-shm`)).rejects.toMatchObject({ code: "ENOENT" });
|
||||
});
|
||||
|
||||
it("previews and forgets sessions without fetching unrelated session bodies", async () => {
|
||||
const db = openOpenClawAgentDatabase({ agentId: "main" }).db;
|
||||
const insert = db.prepare(`INSERT INTO memory_index_chunks
|
||||
|
|
|
|||
|
|
@ -253,7 +253,10 @@ async function forgetWorkspaceMemory(
|
|||
readSessionIngestionState(workspaceDir),
|
||||
readMemoryPreimages(workspaceDir),
|
||||
listMemoryArtifactProvenance({ workspaceDir }),
|
||||
listSessionTranscriptCorpusEntriesForAgent(params.agentId),
|
||||
listSessionTranscriptCorpusEntriesForAgent(params.agentId, {
|
||||
readOnly: true,
|
||||
includeContentRevision: false,
|
||||
}),
|
||||
]);
|
||||
const sessionKeys = new Set(targets.map((target) => target.sessionKey));
|
||||
const curatedWrites = new Map(
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
import type { DatabaseSync } from "node:sqlite";
|
||||
import { DatabaseSync } from "node:sqlite";
|
||||
import {
|
||||
listSessionTranscriptCorpusEntriesForAgent,
|
||||
sessionPathForFile,
|
||||
|
|
@ -120,6 +120,68 @@ describe("memory manager reads", () => {
|
|||
expect(sourceReads.reduce((total, read) => total + read.rows, 0)).toBeLessThanOrEqual(7);
|
||||
});
|
||||
|
||||
it("inspects readonly session corpus diagnostics without changing the published index", async () => {
|
||||
const sessionId = "diagnostic-corpus";
|
||||
await fixture.seedSessionTranscript({
|
||||
sessionId,
|
||||
messages: [
|
||||
{
|
||||
role: "user",
|
||||
timestamp: Date.now(),
|
||||
content: "Violet diagnostic preference remains indexed.",
|
||||
senderIsOwner: true,
|
||||
},
|
||||
],
|
||||
});
|
||||
const cfg = fixture.createConfig({
|
||||
provider: "none",
|
||||
sources: ["sessions"],
|
||||
sessionMemory: true,
|
||||
});
|
||||
const writer = await fixture.getFreshManager(cfg, "cli");
|
||||
await writer.sync({ reason: "diagnostic-baseline", force: true });
|
||||
const database: unknown = Reflect.get(writer, "db");
|
||||
if (!(database instanceof DatabaseSync)) {
|
||||
throw new Error("Expected the fixture's actual published database");
|
||||
}
|
||||
const snapshot = () => ({
|
||||
sources: database.prepare("SELECT * FROM memory_index_sources ORDER BY id").all(),
|
||||
chunks: database.prepare("SELECT * FROM memory_index_chunks ORDER BY id").all(),
|
||||
provenance: database
|
||||
.prepare("SELECT * FROM memory_index_chunk_provenance ORDER BY chunk_id")
|
||||
.all(),
|
||||
metadata: database.prepare("SELECT * FROM memory_index_meta ORDER BY key").all(),
|
||||
revision: database.prepare("SELECT revision FROM memory_index_state WHERE id = 1").get(),
|
||||
nodes: database.prepare("SELECT * FROM session_nodes ORDER BY session_key").all(),
|
||||
windows: database.prepare("SELECT * FROM session_windows ORDER BY session_id").all(),
|
||||
});
|
||||
const before = snapshot();
|
||||
expect(before.sources).toHaveLength(1);
|
||||
expect(before.sources[0]).toMatchObject({
|
||||
path: sessionPathForSessionIdentity("main", sessionId),
|
||||
source: "sessions",
|
||||
});
|
||||
expect(before.chunks).toHaveLength(1);
|
||||
expect(before.chunks[0]).toMatchObject({
|
||||
source: "sessions",
|
||||
text: expect.stringContaining("Violet diagnostic preference remains indexed."),
|
||||
});
|
||||
|
||||
const diagnostic = await fixture.getFreshManager(cfg, "status", true);
|
||||
try {
|
||||
expect(diagnostic.status()).toMatchObject({
|
||||
files: 1,
|
||||
chunks: 1,
|
||||
sourceCounts: [{ source: "sessions", files: 1, chunks: 1, eligible: 1, issues: [] }],
|
||||
});
|
||||
expect(snapshot()).toEqual(before);
|
||||
} finally {
|
||||
await diagnostic.close();
|
||||
}
|
||||
expect(database.isOpen).toBe(true);
|
||||
expect(snapshot()).toEqual(before);
|
||||
});
|
||||
|
||||
it("reuses diagnostic cache totals and the synchronous sync existence check", async () => {
|
||||
const cfg = fixture.createConfig({ provider: "none", cacheEnabled: true });
|
||||
const manager = await fixture.getFreshManager(cfg, "cli");
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
// Session/runtime facade for memory transcript helpers.
|
||||
import path from "node:path";
|
||||
import { isValidAgentId, normalizeAgentId } from "@openclaw/normalization-core/agent-id";
|
||||
import { cloneEnvWithPlatformSemantics } from "../../../../src/config/config-env-vars.js";
|
||||
import {
|
||||
readTranscriptExportSnapshotReadOnlySync,
|
||||
readTranscriptStatsBatchReadOnlySync,
|
||||
|
|
@ -18,7 +19,7 @@ export {
|
|||
} from "../../../../src/config/sessions/session-accessor.js";
|
||||
export { isIncognitoSessionKey } from "../../../../src/routing/session-key.js";
|
||||
export { isIncognitoOpenClawAgentSqlitePath } from "../../../../src/state/openclaw-agent-db.paths.js";
|
||||
export { cloneEnvWithPlatformSemantics } from "../../../../src/config/config-env-vars.js";
|
||||
export { cloneEnvWithPlatformSemantics };
|
||||
|
||||
/** Keep worker launch machinery behind the memory host's existing lazy runtime bridge. */
|
||||
export async function prepareSessionEntryInWorker(
|
||||
|
|
@ -31,6 +32,17 @@ export async function prepareSessionEntryInWorker(
|
|||
return prepare(...args);
|
||||
}
|
||||
|
||||
export async function readSessionEntrySummariesInWorker(
|
||||
input: Parameters<
|
||||
typeof import("../../../../src/config/sessions/session-entry-read-runtime.js").readSessionEntrySummariesInWorker
|
||||
>[0],
|
||||
) {
|
||||
const captured = { ...input, env: cloneEnvWithPlatformSemantics(input.env ?? process.env) };
|
||||
const { readSessionEntrySummariesInWorker: read } =
|
||||
await import("../../../../src/config/sessions/session-entry-read-runtime.js");
|
||||
return read(captured);
|
||||
}
|
||||
|
||||
export { resolveSessionAgentId } from "../../../../src/agents/agent-scope.js";
|
||||
export { stripInternalRuntimeContext } from "../../../../src/agents/internal-runtime-context.js";
|
||||
export { isHeartbeatUserMessage } from "../../../../src/auto-reply/heartbeat-filter.js";
|
||||
|
|
|
|||
|
|
@ -1,17 +1,20 @@
|
|||
import fsSync from "node:fs";
|
||||
import fs from "node:fs/promises";
|
||||
import path from "node:path";
|
||||
import { StatementSync } from "node:sqlite";
|
||||
import {
|
||||
clearConfigCache,
|
||||
clearRuntimeConfigSnapshot,
|
||||
} from "openclaw/plugin-sdk/runtime-config-snapshot";
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import { upsertSessionEntryCore } from "../../../../src/config/sessions/session-accessor.js";
|
||||
import { closeOpenClawAgentDatabasesAsync } from "../../../../src/state/openclaw-agent-db-lifecycle.js";
|
||||
import { registerOpenClawAgentDatabase } from "../../../../src/state/openclaw-agent-db-registry.js";
|
||||
import { getOpenClawAgentDatabaseIfOpen } from "../../../../src/state/openclaw-agent-db.js";
|
||||
import { tableExists } from "../../../../src/state/openclaw-state-db-schema-helpers.js";
|
||||
import { withOpenClawTestState } from "../../../../src/test-utils/openclaw-test-state.js";
|
||||
import { createDeferred } from "../../../../test/helpers/promise.js";
|
||||
import { observeSqliteReadSql } from "../../../../test/helpers/sqlite-statement-execution-counter.js";
|
||||
import { listSessionTranscriptCorpusEntriesForAgent } from "./session-files.js";
|
||||
import { listSessionTranscriptCorpusEntriesForAgentSync } from "./session-transcript-corpus.js";
|
||||
|
||||
|
|
@ -33,12 +36,14 @@ function pauseDirectoryDiscovery(sessionsDir: string) {
|
|||
|
||||
describe("listSessionTranscriptCorpusEntriesForAgent", () => {
|
||||
it.each([
|
||||
{ includeContentRevision: true, archiveTablePresent: true },
|
||||
{ includeContentRevision: false, archiveTablePresent: true },
|
||||
{ includeContentRevision: false, archiveTablePresent: false },
|
||||
{ includeContentRevision: true, archiveTablePresent: true, readOnly: false },
|
||||
{ includeContentRevision: false, archiveTablePresent: true, readOnly: false },
|
||||
{ includeContentRevision: false, archiveTablePresent: false, readOnly: false },
|
||||
{ includeContentRevision: false, archiveTablePresent: true, readOnly: true },
|
||||
{ includeContentRevision: false, archiveTablePresent: false, readOnly: true },
|
||||
])(
|
||||
"preserves corpus selection with content revisions $includeContentRevision and archive table $archiveTablePresent",
|
||||
async ({ includeContentRevision, archiveTablePresent }) => {
|
||||
"preserves corpus selection with content revisions $includeContentRevision and archive table $archiveTablePresent (readOnly: $readOnly)",
|
||||
async ({ includeContentRevision, archiveTablePresent, readOnly }) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
|
||||
const sessionsDir = state.sessionsDir();
|
||||
await fs.mkdir(sessionsDir, { recursive: true });
|
||||
|
|
@ -61,7 +66,7 @@ describe("listSessionTranscriptCorpusEntriesForAgent", () => {
|
|||
}
|
||||
expect(tableExists(db, "session_transcript_archives")).toBe(archiveTablePresent);
|
||||
|
||||
const options = { includeContentRevision };
|
||||
const options = { includeContentRevision, readOnly };
|
||||
const expected = listSessionTranscriptCorpusEntriesForAgentSync("main", options);
|
||||
const actual = await listSessionTranscriptCorpusEntriesForAgent("main", options);
|
||||
|
||||
|
|
@ -80,9 +85,13 @@ describe("listSessionTranscriptCorpusEntriesForAgent", () => {
|
|||
},
|
||||
);
|
||||
|
||||
it.each([false, true])(
|
||||
"classifies prompt-rich entries without decoding saved prompts (readOnly: %s)",
|
||||
async (readOnly) => {
|
||||
it.each([
|
||||
{ readOnly: false, includeRetainedSqlite: true },
|
||||
{ readOnly: true, includeRetainedSqlite: true },
|
||||
{ readOnly: true, includeRetainedSqlite: false },
|
||||
])(
|
||||
"classifies prompt-rich entries without decoding saved prompts (readOnly: $readOnly, retained: $includeRetainedSqlite)",
|
||||
async ({ readOnly, includeRetainedSqlite }) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
|
||||
const sessionKey = "agent:main:corpus-metadata";
|
||||
const storePath = path.join(state.sessionsDir(), "sessions.json");
|
||||
|
|
@ -105,10 +114,23 @@ describe("listSessionTranscriptCorpusEntriesForAgent", () => {
|
|||
);
|
||||
|
||||
const parse = vi.spyOn(JSON, "parse");
|
||||
const observed = observeSqliteReadSql(StatementSync.prototype);
|
||||
const isSummaryRead = (sql: string) =>
|
||||
/^\s*select\b[\s\S]*\bfrom\s+["`]?session_nodes\b/i.test(sql) &&
|
||||
/\bentry_json\b|\*/i.test(sql.split(/\bfrom\b/i)[0]!);
|
||||
try {
|
||||
const { db } = getOpenClawAgentDatabaseIfOpen({ agentId: "main", env: state.env })!;
|
||||
expect(
|
||||
db.prepare('select "session_key", "entry_json" from "session_nodes" where 0').all(),
|
||||
).toEqual([]);
|
||||
expect(observed.queries.filter(isSummaryRead)).toHaveLength(1);
|
||||
// A cold reader must perform the real projection, not reuse a warm host snapshot.
|
||||
await closeOpenClawAgentDatabasesAsync(state.stateDir);
|
||||
observed.queries.length = 0;
|
||||
parse.mockClear();
|
||||
const entries = await listSessionTranscriptCorpusEntriesForAgent("main", {
|
||||
includeContentRevision: false,
|
||||
includeRetainedSqlite: true,
|
||||
includeRetainedSqlite,
|
||||
readOnly,
|
||||
});
|
||||
expect(entries).toEqual([
|
||||
|
|
@ -124,6 +146,9 @@ describe("listSessionTranscriptCorpusEntriesForAgent", () => {
|
|||
updatedAtMs: persisted?.updatedAt,
|
||||
},
|
||||
]);
|
||||
if (readOnly && !includeRetainedSqlite) {
|
||||
expect(observed.queries.filter(isSummaryRead)).toEqual([]);
|
||||
}
|
||||
const decodedEntries = parse.mock.calls.filter(([json]) =>
|
||||
json.includes('"sessionId":"corpus-metadata"'),
|
||||
);
|
||||
|
|
@ -135,82 +160,153 @@ describe("listSessionTranscriptCorpusEntriesForAgent", () => {
|
|||
).toBe(true);
|
||||
} finally {
|
||||
parse.mockRestore();
|
||||
observed.restore();
|
||||
}
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
it("yields during filesystem discovery while retaining its store and reading fresh entries", async () => {
|
||||
it("preserves admitted metadata parsing when readonly corpus discovery crosses the worker", async () => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
|
||||
const sessionsDir = state.statePath("custom-sessions");
|
||||
const storePath = path.join(sessionsDir, "sessions.json");
|
||||
await fs.mkdir(sessionsDir, { recursive: true });
|
||||
await state.writeConfig({ session: { store: storePath } });
|
||||
clearRuntimeConfigSnapshot();
|
||||
clearConfigCache();
|
||||
await upsertSessionEntryCore(
|
||||
{ sessionKey: "agent:main:chat:original", storePath },
|
||||
{ sessionId: "original", updatedAt: 1 },
|
||||
const sessionKey = "agent:main:admitted-corpus";
|
||||
const sessionId = "admitted-corpus";
|
||||
const storePath = path.join(state.sessionsDir(), "sessions.json");
|
||||
const persisted = await upsertSessionEntryCore(
|
||||
{ sessionKey, storePath },
|
||||
{ sessionId, updatedAt: 10 },
|
||||
);
|
||||
const { entered, release, realpathSpy } = pauseDirectoryDiscovery(sessionsDir);
|
||||
const listing = listSessionTranscriptCorpusEntriesForAgent("main");
|
||||
try {
|
||||
await expect(
|
||||
Promise.race([entered.promise.then(() => "waiting"), listing.then(() => "completed")]),
|
||||
).resolves.toBe("waiting");
|
||||
await new Promise<void>((resolve) => {
|
||||
setImmediate(resolve);
|
||||
});
|
||||
await state.writeConfig({
|
||||
session: { store: state.statePath("replacement", "sessions.json") },
|
||||
});
|
||||
const { db } = getOpenClawAgentDatabaseIfOpen({ agentId: "main", env: state.env })!;
|
||||
const options = { readOnly: true, includeContentRevision: false };
|
||||
const expected = {
|
||||
agentId: "main",
|
||||
artifactKind: "active-session",
|
||||
sessionFile: sessionKey,
|
||||
sessionId,
|
||||
sessionKey,
|
||||
sessionKind: "heartbeat",
|
||||
storePath,
|
||||
transcriptSource: "sqlite",
|
||||
updatedAtMs: persisted?.updatedAt,
|
||||
};
|
||||
expect(listSessionTranscriptCorpusEntriesForAgentSync("main", options)).toEqual([
|
||||
{ ...expected, sessionKind: "interactive" },
|
||||
]);
|
||||
|
||||
// Raw metadata edits preserve the original admitted reader's parsing contract.
|
||||
db.prepare(
|
||||
"UPDATE session_nodes SET entry_json = json_set(entry_json, '$.heartbeatIsolatedBaseSessionKey', ?) WHERE session_key = ?",
|
||||
).run("agent:main:main", sessionKey);
|
||||
const validity = () =>
|
||||
db.prepare("SELECT entry_valid FROM session_nodes WHERE session_key = ?").get(sessionKey)
|
||||
?.entry_valid;
|
||||
expect(validity()).toBe(0);
|
||||
expect(listSessionTranscriptCorpusEntriesForAgentSync("main", options)).toEqual([expected]);
|
||||
expect(validity()).toBe(0);
|
||||
|
||||
await expect(listSessionTranscriptCorpusEntriesForAgent("main", options)).resolves.toEqual([
|
||||
expected,
|
||||
]);
|
||||
expect(validity()).toBe(0);
|
||||
expect(db.isOpen).toBe(true);
|
||||
expect(getOpenClawAgentDatabaseIfOpen({ agentId: "main", env: state.env })?.db === db).toBe(
|
||||
true,
|
||||
);
|
||||
await closeOpenClawAgentDatabasesAsync(state.stateDir);
|
||||
expect(db.isOpen).toBe(false);
|
||||
await expect(listSessionTranscriptCorpusEntriesForAgent("main", options)).rejects.toThrow(
|
||||
"openclaw doctor --fix",
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
it.each([false, true])(
|
||||
"yields during filesystem discovery while retaining its store and reading fresh entries (readOnly: %s)",
|
||||
async (readOnly) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
|
||||
const sessionsDir = state.statePath("custom-sessions");
|
||||
const storePath = path.join(sessionsDir, "sessions.json");
|
||||
await fs.mkdir(sessionsDir, { recursive: true });
|
||||
await state.writeConfig({ session: { store: storePath } });
|
||||
clearRuntimeConfigSnapshot();
|
||||
clearConfigCache();
|
||||
await upsertSessionEntryCore(
|
||||
{ sessionKey: "agent:main:chat:fresh", storePath },
|
||||
{ sessionId: "fresh", updatedAt: 2 },
|
||||
{ sessionKey: "agent:main:chat:original", storePath },
|
||||
{ sessionId: "original", updatedAt: 1 },
|
||||
);
|
||||
release.resolve();
|
||||
|
||||
const entries = await listing;
|
||||
expect(entries.map((entry) => entry.sessionId).toSorted()).toEqual(["fresh", "original"]);
|
||||
expect(entries.every((entry) => entry.storePath === storePath)).toBe(true);
|
||||
} finally {
|
||||
release.resolve();
|
||||
await listing.finally(() => realpathSpy.mockRestore());
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
it("retains the custom store's registry environment while discovery is pending", async () => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
|
||||
const sessionsDir = state.statePath("custom-sessions");
|
||||
const storePath = path.join(sessionsDir, "sessions.json");
|
||||
const unsuffixedDatabase = path.join(sessionsDir, "openclaw-agent.sqlite");
|
||||
await fs.mkdir(sessionsDir, { recursive: true });
|
||||
await state.writeConfig({ session: { store: storePath } });
|
||||
clearRuntimeConfigSnapshot();
|
||||
clearConfigCache();
|
||||
registerOpenClawAgentDatabase({ agentId: "ops", path: unsuffixedDatabase });
|
||||
const { entered, release, realpathSpy } = pauseDirectoryDiscovery(sessionsDir);
|
||||
const listing = listSessionTranscriptCorpusEntriesForAgent("main");
|
||||
try {
|
||||
await expect(
|
||||
Promise.race([entered.promise.then(() => "waiting"), listing.then(() => "completed")]),
|
||||
).resolves.toBe("waiting");
|
||||
Reflect.set(process.env, "OPENCLAW_STATE_DIR", state.statePath("replacement-state"));
|
||||
release.resolve();
|
||||
|
||||
await expect(listing).resolves.toEqual([]);
|
||||
expect(fsSync.existsSync(path.join(sessionsDir, "openclaw-agent.main.sqlite"))).toBe(true);
|
||||
expect(fsSync.existsSync(unsuffixedDatabase)).toBe(false);
|
||||
} finally {
|
||||
release.resolve();
|
||||
await listing.finally(() => {
|
||||
Reflect.set(process.env, "OPENCLAW_STATE_DIR", state.stateDir);
|
||||
realpathSpy.mockRestore();
|
||||
const { entered, release, realpathSpy } = pauseDirectoryDiscovery(sessionsDir);
|
||||
const listing = listSessionTranscriptCorpusEntriesForAgent("main", {
|
||||
readOnly,
|
||||
...(readOnly ? { includeContentRevision: false } : {}),
|
||||
});
|
||||
}
|
||||
});
|
||||
});
|
||||
try {
|
||||
await expect(
|
||||
Promise.race([entered.promise.then(() => "waiting"), listing.then(() => "completed")]),
|
||||
).resolves.toBe("waiting");
|
||||
await new Promise<void>((resolve) => {
|
||||
setImmediate(resolve);
|
||||
});
|
||||
await state.writeConfig({
|
||||
session: { store: state.statePath("replacement", "sessions.json") },
|
||||
});
|
||||
clearRuntimeConfigSnapshot();
|
||||
clearConfigCache();
|
||||
await upsertSessionEntryCore(
|
||||
{ sessionKey: "agent:main:chat:fresh", storePath },
|
||||
{ sessionId: "fresh", updatedAt: 2 },
|
||||
);
|
||||
release.resolve();
|
||||
|
||||
const entries = await listing;
|
||||
expect(entries.map((entry) => entry.sessionId).toSorted()).toEqual(["fresh", "original"]);
|
||||
expect(entries.every((entry) => entry.storePath === storePath)).toBe(true);
|
||||
} finally {
|
||||
release.resolve();
|
||||
await listing.finally(() => realpathSpy.mockRestore());
|
||||
}
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
it.each([false, true])(
|
||||
"retains the custom store's registry environment while discovery is pending (readOnly: %s)",
|
||||
async (readOnly) => {
|
||||
await withOpenClawTestState({ scenario: "minimal" }, async (state) => {
|
||||
const sessionsDir = state.statePath("custom-sessions");
|
||||
const storePath = path.join(sessionsDir, "sessions.json");
|
||||
const unsuffixedDatabase = path.join(sessionsDir, "openclaw-agent.sqlite");
|
||||
await fs.mkdir(sessionsDir, { recursive: true });
|
||||
await state.writeConfig({ session: { store: storePath } });
|
||||
clearRuntimeConfigSnapshot();
|
||||
clearConfigCache();
|
||||
registerOpenClawAgentDatabase({ agentId: "ops", path: unsuffixedDatabase });
|
||||
const { entered, release, realpathSpy } = pauseDirectoryDiscovery(sessionsDir);
|
||||
const listing = listSessionTranscriptCorpusEntriesForAgent("main", {
|
||||
readOnly,
|
||||
...(readOnly ? { includeContentRevision: false } : {}),
|
||||
});
|
||||
try {
|
||||
await expect(
|
||||
Promise.race([entered.promise.then(() => "waiting"), listing.then(() => "completed")]),
|
||||
).resolves.toBe("waiting");
|
||||
Reflect.set(process.env, "OPENCLAW_STATE_DIR", state.statePath("replacement-state"));
|
||||
release.resolve();
|
||||
|
||||
await expect(listing).resolves.toEqual([]);
|
||||
expect(fsSync.existsSync(path.join(sessionsDir, "openclaw-agent.main.sqlite"))).toBe(
|
||||
!readOnly,
|
||||
);
|
||||
expect(fsSync.existsSync(unsuffixedDatabase)).toBe(false);
|
||||
if (readOnly) {
|
||||
expect(fsSync.existsSync(state.statePath("replacement-state"))).toBe(false);
|
||||
}
|
||||
} finally {
|
||||
release.resolve();
|
||||
await listing.finally(() => {
|
||||
Reflect.set(process.env, "OPENCLAW_STATE_DIR", state.stateDir);
|
||||
realpathSpy.mockRestore();
|
||||
});
|
||||
}
|
||||
});
|
||||
},
|
||||
);
|
||||
});
|
||||
|
|
|
|||
|
|
@ -18,6 +18,7 @@ import {
|
|||
listSessionTranscriptArchivesReadOnly,
|
||||
listSessionTranscriptInstances,
|
||||
parseUsageCountedSessionIdFromFileName,
|
||||
readSessionEntrySummariesInWorker,
|
||||
readTranscriptContentRevisionSync,
|
||||
resolveSessionAgentId,
|
||||
resolveSessionTranscriptsDirForAgent,
|
||||
|
|
@ -384,20 +385,12 @@ function projectSessionTranscriptCorpusEntries(
|
|||
scope: ReturnType<typeof resolveSessionTranscriptCorpusScope>,
|
||||
options: SessionTranscriptCorpusOptions,
|
||||
artifacts: readonly SessionTranscriptCorpusArtifact[],
|
||||
sessionEntries: readonly SessionEntrySummary[],
|
||||
): SessionTranscriptCorpusEntry[] {
|
||||
const { cfg, env, normalizedAgentId, storePath, isSharedFixedStore } = scope;
|
||||
const includeContentRevision = options.includeContentRevision !== false;
|
||||
const activeEntriesBySessionId = new Map<string, SessionTranscriptCorpusEntry>();
|
||||
const entryOwnersBySessionId = new Map<string, string>();
|
||||
const listEntries =
|
||||
options.readOnly === true ? listSessionEntriesReadOnly : listSessionEntriesCore;
|
||||
const sessionEntries = listEntries({
|
||||
agentId: normalizedAgentId,
|
||||
env,
|
||||
hydrateSkillPromptRefs: false,
|
||||
projection: "list",
|
||||
storePath,
|
||||
});
|
||||
const retainedInstances = options.includeRetainedSqlite
|
||||
? listSessionTranscriptInstances({
|
||||
agentId: normalizedAgentId,
|
||||
|
|
@ -535,7 +528,27 @@ export function listSessionTranscriptCorpusEntriesForAgentSync(
|
|||
}
|
||||
}
|
||||
}
|
||||
return projectSessionTranscriptCorpusEntries(scope, options, artifacts);
|
||||
return projectSessionTranscriptCorpusEntries(
|
||||
scope,
|
||||
options,
|
||||
artifacts,
|
||||
readCorpusSessionEntries(scope, options),
|
||||
);
|
||||
}
|
||||
|
||||
function readCorpusSessionEntries(
|
||||
scope: ReturnType<typeof resolveSessionTranscriptCorpusScope>,
|
||||
options: SessionTranscriptCorpusOptions,
|
||||
): SessionEntrySummary[] {
|
||||
const listEntries =
|
||||
options.readOnly === true ? listSessionEntriesReadOnly : listSessionEntriesCore;
|
||||
return listEntries({
|
||||
agentId: scope.normalizedAgentId,
|
||||
env: scope.env,
|
||||
hydrateSkillPromptRefs: false,
|
||||
projection: "list",
|
||||
storePath: scope.storePath,
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -589,5 +602,13 @@ export async function listSessionTranscriptCorpusEntriesForAgent(
|
|||
}
|
||||
// Read current session ownership only after the filesystem awaits, while retaining
|
||||
// the caller's resolved store and alias configuration for this complete projection.
|
||||
return projectSessionTranscriptCorpusEntries(scope, capturedOptions, artifacts);
|
||||
const sessionEntries =
|
||||
capturedOptions.readOnly === true
|
||||
? await readSessionEntrySummariesInWorker({
|
||||
agentId: scope.normalizedAgentId,
|
||||
env: scope.env,
|
||||
storePath: scope.storePath,
|
||||
})
|
||||
: readCorpusSessionEntries(scope, capturedOptions);
|
||||
return projectSessionTranscriptCorpusEntries(scope, capturedOptions, artifacts, sessionEntries);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -104,13 +104,18 @@ export function readSelectedSessionEntriesInDatabase(
|
|||
*/
|
||||
export function listSessionEntriesReadOnly(
|
||||
scope: SessionEntryListScope = {},
|
||||
options: { deferParticipants?: true } = {},
|
||||
options: {
|
||||
deferParticipants?: true;
|
||||
continuation?: CanonicalSessionReaderContinuation;
|
||||
} = {},
|
||||
): SessionEntrySummary[] {
|
||||
const resolved = resolveSqliteScope({ ...scope, sessionKey: "" });
|
||||
const result = withOpenClawAgentDatabaseReadOnly(
|
||||
(database) => listSqliteSessionEntriesFromDatabase(database, resolved, scope, options),
|
||||
toDatabaseOptions(resolved),
|
||||
);
|
||||
const result = withOpenClawAgentDatabaseReadOnly((database) => {
|
||||
const read = () => listSqliteSessionEntriesFromDatabase(database, resolved, scope, options);
|
||||
return options.continuation
|
||||
? readWithCanonicalSessionReaderContinuation(database, options.continuation, read)
|
||||
: read();
|
||||
}, toDatabaseOptions(resolved));
|
||||
return result.found ? result.value : [];
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ import { captureOpenClawAgentDatabaseExecution } from "../../state/openclaw-agen
|
|||
import { runOpenClawAgentWorkerWrite } from "../../state/openclaw-agent-write-admission.js";
|
||||
import { cloneEnvWithPlatformSemantics } from "../config-env-vars.js";
|
||||
import { resolveStateDir } from "../state-dir.js";
|
||||
import { listSessionEntriesReadOnly } from "./session-accessor.sqlite-entry-list.read.js";
|
||||
import {
|
||||
loadSessionEntry,
|
||||
loadSessionEntryReadOnlyResultInScope,
|
||||
|
|
@ -412,6 +413,38 @@ type SessionStoreWorkerReadScope = {
|
|||
env?: NodeJS.ProcessEnv;
|
||||
};
|
||||
|
||||
/** Read descriptive summaries through the original store selection and reader lifetime. */
|
||||
export async function readSessionEntrySummariesInWorker(input: SessionStoreWorkerReadScope) {
|
||||
const { scope, agentId } = captureSessionEntryReadScope({ ...input, sessionKey: "" });
|
||||
if (isNativeSessionEntryRead(scope, agentId)) {
|
||||
// Process-held transcripts keep their existing native reader until its worker cutover.
|
||||
return listSessionEntriesReadOnly({
|
||||
...scope,
|
||||
projection: "list",
|
||||
hydrateSkillPromptRefs: false,
|
||||
});
|
||||
}
|
||||
return withSessionStoreReaderInWorker(
|
||||
{ ...input, env: scope.env, storePath: scope.storePath ?? input.storePath },
|
||||
async (owner, database, continuation, assertCurrent) => {
|
||||
assertCurrent();
|
||||
const entries = await owner.readEntries(
|
||||
{
|
||||
agentId: database.agentId,
|
||||
storePath: database.path,
|
||||
env: database.env,
|
||||
projection: "list",
|
||||
hydrateSkillPromptRefs: false,
|
||||
},
|
||||
continuation,
|
||||
);
|
||||
assertCurrent();
|
||||
return entries;
|
||||
},
|
||||
{ backing: true, dataOnly: true },
|
||||
);
|
||||
}
|
||||
|
||||
type SessionEntryWorkerRead = SessionStoreWorkerReadScope &
|
||||
SessionExactEntriesWorkerSelection & {
|
||||
lifecycleSessionKey?: string;
|
||||
|
|
|
|||
|
|
@ -303,12 +303,15 @@ export function createSessionHistoryWorkerReaders(
|
|||
(input) => ({ kind: "session-diagnostic-text", ...input }),
|
||||
(value) => value.text,
|
||||
),
|
||||
readEntries: reader(
|
||||
"session-entry-list",
|
||||
"entries",
|
||||
(scope) => ({ kind: "session-entry-list", scope }),
|
||||
(value) => value.entries,
|
||||
),
|
||||
readEntries: async (scope, continuation) =>
|
||||
runRequest(
|
||||
() => ({ kind: "session-entry-list", scope, continuation }),
|
||||
JSON.stringify({ scope, continuation }).length * 2,
|
||||
(value) => {
|
||||
assertResultKind(value, "session-entry-list", "entries");
|
||||
return value.entries;
|
||||
},
|
||||
),
|
||||
readIdentityEvidence: reader(
|
||||
"session-identity-evidence",
|
||||
"identity evidence",
|
||||
|
|
|
|||
|
|
@ -332,6 +332,7 @@ type SessionEntryListWorkerInput = {
|
|||
kind: "session-entry-list";
|
||||
database: { agentId: string; path: string };
|
||||
scope: SessionEntryListScope;
|
||||
continuation?: CanonicalSessionReaderContinuation;
|
||||
};
|
||||
|
||||
export type SessionExactEntriesWorkerInput = {
|
||||
|
|
@ -651,7 +652,10 @@ export type SessionHistoryWorkerDatabase = {
|
|||
signal?: AbortSignal,
|
||||
) => Promise<SessionExactEntriesWorkerResult>;
|
||||
readRowFacts: SessionHistoryReader<SessionRowFactsWorkerInput>;
|
||||
readEntries: (scope: SessionEntryListWorkerInput["scope"]) => Promise<SessionEntrySummary[]>;
|
||||
readEntries: (
|
||||
scope: SessionEntryListWorkerInput["scope"],
|
||||
continuation?: CanonicalSessionReaderContinuation,
|
||||
) => Promise<SessionEntrySummary[]>;
|
||||
readEntryResult: SessionHistoryReader<
|
||||
SessionEntryReadWorkerInput,
|
||||
import("@openclaw/normalization-core/result").Result<
|
||||
|
|
|
|||
|
|
@ -358,10 +358,13 @@ serveOwnedWorkerTasks(
|
|||
await import("./session-accessor.sqlite-entry-list.read.js");
|
||||
return {
|
||||
kind: "session-entry-list" as const,
|
||||
entries: listSessionEntriesReadOnly({
|
||||
...request.scope,
|
||||
env: cloneEnvWithPlatformSemantics(request.scope.env ?? process.env),
|
||||
}),
|
||||
entries: listSessionEntriesReadOnly(
|
||||
{
|
||||
...request.scope,
|
||||
env: cloneEnvWithPlatformSemantics(request.scope.env ?? process.env),
|
||||
},
|
||||
{ continuation: request.continuation },
|
||||
),
|
||||
};
|
||||
}
|
||||
if (request.kind === "usage-cache") {
|
||||
|
|
|
|||
|
|
@ -549,6 +549,7 @@ export const databaseWorkerCoreTestFiles = [
|
|||
"test/runtime-agent.codex-initialization.integration.test.ts",
|
||||
"packages/memory-host-sdk/src/host/session-memory-sync.test.ts",
|
||||
"packages/memory-host-sdk/src/host/session-files-archive-identity.test.ts",
|
||||
"packages/memory-host-sdk/src/host/session-transcript-corpus.test.ts",
|
||||
"src/agents/harness/native-hook-relay-store.test.ts",
|
||||
"src/agents/harness/native-hook-relay.approval-binding.test.ts",
|
||||
"src/agents/harness/native-hook-relay.approval-wait.test.ts",
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue