mirror of
https://github.com/openclaw/openclaw.git
synced 2026-09-28 05:54:09 +08:00
perf(tasks): batch session access reads within scan slices (#137811)
Resolve only the requested canonical session candidates through one read-only admission per synchronous registry slice, and refresh selected entries for final authorization. Preserve scalar lookup ordering, per-target failures, cursor/revision fences, and the original final-batch fairness repair. Avoid repeated cold store validation without materializing unrelated warm session inventories. No configuration, schema, or public RPC changes.
This commit is contained in:
@@ -26,6 +26,7 @@ import {
|
||||
listSessionEntryKeysReadOnly,
|
||||
loadExactSessionEntry,
|
||||
loadExactSessionEntryCandidates,
|
||||
loadExactSessionEntryCandidatesReadOnlyBatch,
|
||||
loadExactSessionEntryReadOnly,
|
||||
loadSessionEntry,
|
||||
loadSessionEntryReadOnly,
|
||||
@@ -80,6 +81,7 @@ export {
|
||||
listSessionEntryKeysReadOnly,
|
||||
loadExactSessionEntry,
|
||||
loadExactSessionEntryCandidates,
|
||||
loadExactSessionEntryCandidatesReadOnlyBatch,
|
||||
loadExactSessionEntryReadOnly,
|
||||
loadSessionEntry,
|
||||
loadSessionEntryReadOnly,
|
||||
|
||||
@@ -22,6 +22,7 @@ import {
|
||||
hasSessionEntriesByStatusReadOnly,
|
||||
listSessionEntriesCore,
|
||||
listSessionEntriesReadOnly,
|
||||
loadExactSessionEntryCandidatesReadOnlyBatch,
|
||||
loadExactSessionEntryReadOnly,
|
||||
openSessionEntryReadView,
|
||||
readSessionTranscriptTitleProbeBatch,
|
||||
@@ -101,6 +102,11 @@ describe("session accessor readonly listing", () => {
|
||||
|
||||
expect(listSessionEntriesReadOnly({ agentId, env })).toEqual([]);
|
||||
expect(openSessionEntryReadView({ agentId, env }).entries()).toEqual([]);
|
||||
expect(
|
||||
loadExactSessionEntryCandidatesReadOnlyBatch([
|
||||
{ agentId, env, sessionKeys: [`agent:${agentId}:main`] },
|
||||
]),
|
||||
).toEqual([{ ok: true, value: [] }]);
|
||||
expect(
|
||||
readSessionStoreSummaryReadOnly(
|
||||
{ agentId, env },
|
||||
@@ -194,13 +200,38 @@ describe("session accessor readonly listing", () => {
|
||||
});
|
||||
|
||||
const retainedScope = { ...scope, sessionKey: "agent:main:retained" };
|
||||
const exactReadFailure = {
|
||||
ok: false,
|
||||
error: expect.objectContaining({ message: expect.stringContaining("openclaw doctor --fix") }),
|
||||
};
|
||||
for (const projection of ["full", "list"] as const) {
|
||||
expect(loadExactSessionEntryReadOnly({ ...retainedScope, projection })).toBeUndefined();
|
||||
for (const key of ["bad-json", "bad-timestamp"]) {
|
||||
for (const key of ["pending", "bad-json", "bad-timestamp"]) {
|
||||
expect(() =>
|
||||
loadExactSessionEntryReadOnly({ ...scope, sessionKey: `agent:main:${key}`, projection }),
|
||||
).toThrow("openclaw doctor --fix");
|
||||
}
|
||||
const grouped = loadExactSessionEntryCandidatesReadOnlyBatch(
|
||||
[
|
||||
["agent:main:tie-b"],
|
||||
["agent:main:pending"],
|
||||
["agent:main:bad-json"],
|
||||
["agent:main:tie-a", "agent:main:bad-timestamp"],
|
||||
["agent:main:tie-a"],
|
||||
[retainedScope.sessionKey, "agent:main:missing"],
|
||||
].map((sessionKeys) => ({ agentId: scope.agentId, env, sessionKeys, projection })),
|
||||
);
|
||||
expect(grouped).toMatchObject([
|
||||
{
|
||||
ok: true,
|
||||
value: [{ sessionKey: "agent:main:tie-b", entry: { sessionId: "tie-b" } }],
|
||||
},
|
||||
exactReadFailure,
|
||||
exactReadFailure,
|
||||
exactReadFailure,
|
||||
{ ok: true, value: [{ sessionKey: "agent:main:tie-a", entry: { sessionId: "tie-a" } }] },
|
||||
{ ok: true, value: [] },
|
||||
]);
|
||||
}
|
||||
update.run(
|
||||
JSON.stringify({ skillsSnapshot: { prompt: "invalid prompt-only row" } }),
|
||||
@@ -213,6 +244,15 @@ describe("session accessor readonly listing", () => {
|
||||
|
||||
closeOpenClawAgentDatabasesForTest();
|
||||
expect(() => readSessionStoreSummaryReadOnly(scope, options)).toThrow("openclaw doctor --fix");
|
||||
expect(
|
||||
loadExactSessionEntryCandidatesReadOnlyBatch(
|
||||
["agent:main:pending", "agent:main:tie-a"].map((sessionKey) => ({
|
||||
agentId: scope.agentId,
|
||||
env,
|
||||
sessionKeys: [sessionKey],
|
||||
})),
|
||||
),
|
||||
).toEqual([exactReadFailure, exactReadFailure]);
|
||||
});
|
||||
|
||||
it("surfaces missing canonical transcript tables through single and batched reads", async () => {
|
||||
|
||||
@@ -17,7 +17,6 @@ import type { DeliveryContext } from "../../utils/delivery-context.types.js";
|
||||
import { isInternalSessionEffectsKey } from "./internal-session-key.js";
|
||||
import { deriveLastRoutePatch, deriveSessionMetaPatch } from "./metadata.js";
|
||||
import type {
|
||||
ExactSessionEntry,
|
||||
SessionAccessScope,
|
||||
SessionEntryPatchContext,
|
||||
SessionEntryPatchOptions,
|
||||
@@ -77,6 +76,12 @@ import type { GroupKeyResolution, InternalSessionEntry as SessionEntry } from ".
|
||||
import { mergeSessionEntry, mergeSessionEntryPreserveActivity } from "./types.js";
|
||||
|
||||
export { ensureSessionEntrySync } from "./session-accessor.sqlite-initial-entry.js";
|
||||
export {
|
||||
loadExactSessionEntry,
|
||||
loadExactSessionEntryCandidates,
|
||||
loadExactSessionEntryCandidatesReadOnlyBatch,
|
||||
loadExactSessionEntryReadOnly,
|
||||
} from "./session-accessor.sqlite-exact-read.js";
|
||||
|
||||
// Public entry API. Async preparation precedes BEGIN; commit revalidates repository snapshots.
|
||||
|
||||
@@ -133,41 +138,6 @@ export function loadSessionEntryReadOnly(scope: SessionAccessScope): SessionEntr
|
||||
return resolveSessionEntry(scope, { readOnly: true }).existing;
|
||||
}
|
||||
|
||||
/** Loads one exact persisted-key entry from the additive SQLite session store. */
|
||||
export function loadExactSessionEntry(scope: SessionEntryReadScope): ExactSessionEntry | undefined {
|
||||
return loadExactSessionEntryCandidates({
|
||||
...scope,
|
||||
sessionKeys: [scope.sessionKey],
|
||||
readOnly: false,
|
||||
})[0];
|
||||
}
|
||||
|
||||
/** Reads exact candidates for one logical session through a single store admission. */
|
||||
export function loadExactSessionEntryCandidates(
|
||||
scope: Omit<SessionEntryReadScope, "sessionKey"> & {
|
||||
sessionKeys: readonly string[];
|
||||
readOnly: boolean;
|
||||
},
|
||||
): ExactSessionEntry[] {
|
||||
const sessionKeys = scope.sessionKeys.map((key) => key.trim()).filter(Boolean);
|
||||
const [sessionKey] = sessionKeys;
|
||||
if (!sessionKey) {
|
||||
return [];
|
||||
}
|
||||
const resolved = resolveSqliteScope({ ...scope, sessionKey });
|
||||
// Alias candidates share a store; fresh handles must not rescan canonical state per key.
|
||||
const read = (database: Pick<OpenClawAgentDatabase, "agentId" | "db">) =>
|
||||
sessionKeys.flatMap((key) => {
|
||||
const entry = readExactSessionEntryRowValidated(database, key, scope.projection)?.entry;
|
||||
return entry ? [{ sessionKey: key, entry }] : [];
|
||||
});
|
||||
if (!scope.readOnly) {
|
||||
return read(openOpenClawAgentDatabase(toDatabaseOptions(resolved)));
|
||||
}
|
||||
const result = withOpenClawAgentDatabaseReadOnly(read, toDatabaseOptions(resolved));
|
||||
return result.found ? result.value : [];
|
||||
}
|
||||
|
||||
/** Lists persisted session keys without materializing their entry JSON. */
|
||||
export function listSessionEntryKeysReadOnly(
|
||||
scope: Partial<Omit<SessionAccessScope, "sessionKey">> = {},
|
||||
@@ -183,17 +153,6 @@ export function listSessionEntryKeysReadOnly(
|
||||
return result.found ? result.value : [];
|
||||
}
|
||||
|
||||
/** Exact persisted-key probe on the read-only handle, for per-row hot paths. */
|
||||
export function loadExactSessionEntryReadOnly(
|
||||
scope: SessionEntryReadScope,
|
||||
): ExactSessionEntry | undefined {
|
||||
return loadExactSessionEntryCandidates({
|
||||
...scope,
|
||||
sessionKeys: [scope.sessionKey],
|
||||
readOnly: true,
|
||||
})[0];
|
||||
}
|
||||
|
||||
/** Lists direct child rows without cloning or rebuilding the complete session store. */
|
||||
export function listSessionChildEntriesReadOnly(
|
||||
scope: SessionEntryReadScope,
|
||||
|
||||
@@ -0,0 +1,144 @@
|
||||
import { err, ok, type Result } from "@openclaw/normalization-core/result";
|
||||
import { withOpenClawAgentDatabaseReadOnly } from "../../state/openclaw-agent-db-readonly.js";
|
||||
import {
|
||||
openOpenClawAgentDatabase,
|
||||
resolveOpenClawAgentSqlitePath,
|
||||
type OpenClawAgentDatabase,
|
||||
type OpenClawAgentDatabaseOptions,
|
||||
} from "../../state/openclaw-agent-db.js";
|
||||
import type { ExactSessionEntry } from "./session-accessor.sqlite-contract.js";
|
||||
import { readExactSessionEntryRowValidated } from "./session-accessor.sqlite-entry-store.js";
|
||||
import { resolveSqliteScope, toDatabaseOptions } from "./session-accessor.sqlite-scope.js";
|
||||
import type { SessionEntryReadScope } from "./session-accessor.types.js";
|
||||
import { assertCanonicalSqliteSessionKeysCurrent } from "./session-canonical-key.js";
|
||||
|
||||
/** Loads one exact persisted-key entry from the additive SQLite session store. */
|
||||
export function loadExactSessionEntry(scope: SessionEntryReadScope): ExactSessionEntry | undefined {
|
||||
return loadExactSessionEntryCandidates({
|
||||
...scope,
|
||||
sessionKeys: [scope.sessionKey],
|
||||
readOnly: false,
|
||||
})[0];
|
||||
}
|
||||
|
||||
/** Reads exact candidates for one logical session through a single store admission. */
|
||||
export function loadExactSessionEntryCandidates(
|
||||
scope: Omit<SessionEntryReadScope, "sessionKey"> & {
|
||||
sessionKeys: readonly string[];
|
||||
readOnly: boolean;
|
||||
},
|
||||
): ExactSessionEntry[] {
|
||||
const sessionKeys = scope.sessionKeys.map((key) => key.trim()).filter(Boolean);
|
||||
const [sessionKey] = sessionKeys;
|
||||
if (!sessionKey) {
|
||||
return [];
|
||||
}
|
||||
const resolved = resolveSqliteScope({ ...scope, sessionKey });
|
||||
// Alias candidates share a store; fresh handles must not rescan canonical state per key.
|
||||
const read = (database: Pick<OpenClawAgentDatabase, "agentId" | "db">) =>
|
||||
sessionKeys.flatMap((key) => {
|
||||
const entry = readExactSessionEntryRowValidated(database, key, scope.projection)?.entry;
|
||||
return entry ? [{ sessionKey: key, entry }] : [];
|
||||
});
|
||||
if (!scope.readOnly) {
|
||||
return read(openOpenClawAgentDatabase(toDatabaseOptions(resolved)));
|
||||
}
|
||||
const result = withOpenClawAgentDatabaseReadOnly(read, toDatabaseOptions(resolved));
|
||||
return result.found ? result.value : [];
|
||||
}
|
||||
|
||||
/** Exact persisted-key probe on the read-only handle, for per-row hot paths. */
|
||||
export function loadExactSessionEntryReadOnly(
|
||||
scope: SessionEntryReadScope,
|
||||
): ExactSessionEntry | undefined {
|
||||
return loadExactSessionEntryCandidates({
|
||||
...scope,
|
||||
sessionKeys: [scope.sessionKey],
|
||||
readOnly: true,
|
||||
})[0];
|
||||
}
|
||||
|
||||
/** Read requested keys through synchronous store/projection groups. */
|
||||
export function loadExactSessionEntryCandidatesReadOnlyBatch(
|
||||
scopes: readonly (Omit<SessionEntryReadScope, "sessionKey"> & {
|
||||
sessionKeys: readonly string[];
|
||||
})[],
|
||||
): Array<Result<ExactSessionEntry[], unknown>> {
|
||||
const results: Array<Result<ExactSessionEntry[], unknown>> = scopes.map(() => ok([]));
|
||||
const groups = new Map<
|
||||
string,
|
||||
{
|
||||
options: OpenClawAgentDatabaseOptions;
|
||||
projection: SessionEntryReadScope["projection"];
|
||||
requests: Array<{ index: number; sessionKeys: string[] }>;
|
||||
}
|
||||
>();
|
||||
for (const [index, scope] of scopes.entries()) {
|
||||
const sessionKeys = scope.sessionKeys.map((key) => key.trim()).filter(Boolean);
|
||||
const [sessionKey] = sessionKeys;
|
||||
if (!sessionKey) {
|
||||
continue;
|
||||
}
|
||||
try {
|
||||
const options = toDatabaseOptions(resolveSqliteScope({ ...scope, sessionKey }));
|
||||
const groupKey = [
|
||||
options.agentId,
|
||||
resolveOpenClawAgentSqlitePath(options),
|
||||
scope.projection ?? "full",
|
||||
].join("\u0000");
|
||||
const group = groups.get(groupKey) ?? { options, projection: scope.projection, requests: [] };
|
||||
group.requests.push({ index, sessionKeys });
|
||||
groups.set(groupKey, group);
|
||||
} catch (error) {
|
||||
results[index] = err(error);
|
||||
}
|
||||
}
|
||||
for (const group of groups.values()) {
|
||||
try {
|
||||
withOpenClawAgentDatabaseReadOnly((database) => {
|
||||
// Admission failures affect this store; an invalid requested row must not
|
||||
// suppress healthy logical targets after a warm handle was validated.
|
||||
assertCanonicalSqliteSessionKeysCurrent(database);
|
||||
const entries = new Map<string, Result<ExactSessionEntry | undefined, unknown>>();
|
||||
const readEntry = (sessionKey: string): Result<ExactSessionEntry | undefined, unknown> => {
|
||||
const cached = entries.get(sessionKey);
|
||||
if (cached) {
|
||||
return cached;
|
||||
}
|
||||
let result: Result<ExactSessionEntry | undefined, unknown>;
|
||||
try {
|
||||
const entry = readExactSessionEntryRowValidated(
|
||||
database,
|
||||
sessionKey,
|
||||
group.projection,
|
||||
)?.entry;
|
||||
result = ok(entry ? { sessionKey, entry } : undefined);
|
||||
} catch (error) {
|
||||
result = err(error);
|
||||
}
|
||||
entries.set(sessionKey, result);
|
||||
return result;
|
||||
};
|
||||
for (const { index, sessionKeys } of group.requests) {
|
||||
const matches: ExactSessionEntry[] = [];
|
||||
results[index] = ok(matches);
|
||||
for (const sessionKey of sessionKeys) {
|
||||
const entry = readEntry(sessionKey);
|
||||
if (!entry.ok) {
|
||||
results[index] = err(entry.error);
|
||||
break;
|
||||
}
|
||||
if (entry.value) {
|
||||
matches.push(entry.value);
|
||||
}
|
||||
}
|
||||
}
|
||||
}, group.options);
|
||||
} catch (error) {
|
||||
for (const { index } of group.requests) {
|
||||
results[index] = err(error);
|
||||
}
|
||||
}
|
||||
}
|
||||
return results;
|
||||
}
|
||||
@@ -153,6 +153,7 @@ export {
|
||||
listSessionEntryKeysReadOnly,
|
||||
loadExactSessionEntry,
|
||||
loadExactSessionEntryCandidates,
|
||||
loadExactSessionEntryCandidatesReadOnlyBatch,
|
||||
loadExactSessionEntryReadOnly,
|
||||
loadSessionEntry,
|
||||
loadSessionEntryReadOnly,
|
||||
|
||||
@@ -0,0 +1,241 @@
|
||||
import { expectDefined } from "@openclaw/normalization-core";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { upsertSessionEntryCore } from "../../config/sessions/session-accessor.js";
|
||||
import { setCanonicalSqliteSessionMainKey } from "../../config/sessions/session-canonical-key.js";
|
||||
import {
|
||||
closeOpenClawAgentDatabasesForTest,
|
||||
listOpenClawAgentDatabasesForTest,
|
||||
openOpenClawAgentDatabase,
|
||||
} from "../../state/openclaw-agent-db.js";
|
||||
import { ensureProfileForEmail } from "../../state/user-profiles.js";
|
||||
import { deleteTaskRecordById } from "../../tasks/runtime-internal.js";
|
||||
import { reloadTaskRegistryFromStore } from "../../tasks/task-registry.js";
|
||||
import { saveTaskRegistryStateToSqlite } from "../../tasks/task-registry.store.sqlite.js";
|
||||
import { resetTaskRegistryForTests } from "../../tasks/task-runtime.test-helpers.js";
|
||||
import { createOpenClawTestState } from "../../test-utils/openclaw-test-state.js";
|
||||
import { readGatewayAccessRevision } from "../gateway-access-revision.js";
|
||||
import { rolePolicyConfig } from "../session-sharing.test-utils.js";
|
||||
import { sessionSharingHandlers } from "./sessions-sharing.js";
|
||||
import {
|
||||
captureRespond,
|
||||
createSnapshotTask,
|
||||
identifiedClient,
|
||||
runTaskHandler,
|
||||
} from "./tasks.test-helpers.js";
|
||||
|
||||
let state: Awaited<ReturnType<typeof createOpenClawTestState>>;
|
||||
beforeEach(async () => {
|
||||
state = await createOpenClawTestState({ scenario: "minimal" });
|
||||
resetTaskRegistryForTests({ persist: false });
|
||||
});
|
||||
afterEach(async () => {
|
||||
resetTaskRegistryForTests({ persist: false });
|
||||
await state.cleanup();
|
||||
});
|
||||
|
||||
describe("task page access snapshots", () => {
|
||||
it.each(["canonical", "main alias", "distinct requesters", "warm"] as const)(
|
||||
"bounds session lookup work across a yielded task page using %s keys",
|
||||
async (mode) => {
|
||||
const sessionKey = "agent:main:cold-requester";
|
||||
const warm = mode === "warm";
|
||||
const requesterKeys =
|
||||
mode === "distinct requesters"
|
||||
? Array.from({ length: 65 }, (_, index) => `agent:main:requester-${index}`)
|
||||
: [sessionKey];
|
||||
const profileId = ensureProfileForEmail("cold-viewer@example.test").id;
|
||||
const config = rolePolicyConfig();
|
||||
if (mode === "main alias") {
|
||||
config.session = { mainKey: "cold-requester" };
|
||||
setCanonicalSqliteSessionMainKey(
|
||||
openOpenClawAgentDatabase({ agentId: "main" }),
|
||||
"cold-requester",
|
||||
);
|
||||
}
|
||||
for (const requesterKey of requesterKeys) {
|
||||
await upsertSessionEntryCore(
|
||||
{ agentId: "main", sessionKey: requesterKey },
|
||||
{ sessionId: `session-${requesterKey}`, updatedAt: 1, visibility: "shared" },
|
||||
);
|
||||
}
|
||||
const unrelatedCount = 24;
|
||||
for (let index = 0; index < unrelatedCount; index += 1) {
|
||||
await upsertSessionEntryCore(
|
||||
{ agentId: "main", sessionKey: `agent:main:cold-unrelated-${index}` },
|
||||
{ sessionId: `cold-unrelated-payload-${index}`, updatedAt: 1 },
|
||||
);
|
||||
}
|
||||
const tasks = Array.from({ length: 65 }, (_, index) =>
|
||||
createSnapshotTask({
|
||||
taskId: `cold-task-${index}`,
|
||||
requesterSessionKey:
|
||||
mode === "main alias" ? "main" : (requesterKeys[index] ?? sessionKey),
|
||||
requesterAgentId: "main",
|
||||
ownerKey: sessionKey,
|
||||
lastEventAt: 2_000 + index,
|
||||
}),
|
||||
);
|
||||
saveTaskRegistryStateToSqlite({
|
||||
tasks: new Map(tasks.map((task) => [task.taskId, task])),
|
||||
deliveryStates: new Map(),
|
||||
});
|
||||
reloadTaskRegistryFromStore();
|
||||
if (!warm) {
|
||||
closeOpenClawAgentDatabasesForTest();
|
||||
}
|
||||
const expectedHandles = listOpenClawAgentDatabasesForTest().length;
|
||||
expect(expectedHandles > 0).toBe(warm);
|
||||
let materializedUnrelated = 0;
|
||||
const materialize = Object.fromEntries;
|
||||
const materializeSpy = vi.spyOn(Object, "fromEntries").mockImplementation((entries) => {
|
||||
const result = materialize(entries);
|
||||
for (const entry of Object.values(result)) {
|
||||
if (
|
||||
entry &&
|
||||
typeof entry === "object" &&
|
||||
"sessionId" in entry &&
|
||||
typeof entry.sessionId === "string" &&
|
||||
entry.sessionId.startsWith("cold-unrelated-payload-")
|
||||
) {
|
||||
materializedUnrelated += 1;
|
||||
}
|
||||
}
|
||||
return result;
|
||||
});
|
||||
let unrelatedParses = 0;
|
||||
const parse = JSON.parse;
|
||||
const parseSpy = vi.spyOn(JSON, "parse").mockImplementation((value, reviver) => {
|
||||
if (value.includes("cold-unrelated-payload-")) {
|
||||
unrelatedParses += 1;
|
||||
}
|
||||
return parse(value, reviver);
|
||||
});
|
||||
const yielded = new Promise<{ parses: number; handles: number }>((resolve) => {
|
||||
setImmediate(() =>
|
||||
resolve({
|
||||
parses: unrelatedParses,
|
||||
handles: listOpenClawAgentDatabasesForTest().length,
|
||||
}),
|
||||
);
|
||||
});
|
||||
try {
|
||||
const { calls, payload } = await runTaskHandler(
|
||||
"tasks.list",
|
||||
{ limit: 100 },
|
||||
config,
|
||||
identifiedClient(["operator.read"], profileId),
|
||||
);
|
||||
expect(calls[0]?.[0]).toBe(true);
|
||||
expect(payload?.tasks?.map((task) => task.id)).toEqual(
|
||||
tasks.toReversed().map((task) => task.taskId),
|
||||
);
|
||||
const slice = await yielded;
|
||||
expect(slice.handles).toBe(expectedHandles);
|
||||
if (!warm) {
|
||||
expect(slice.parses).toBeGreaterThan(0);
|
||||
}
|
||||
expect(listOpenClawAgentDatabasesForTest()).toHaveLength(expectedHandles);
|
||||
expect(materializedUnrelated).toBe(0);
|
||||
// Three synchronous slices plus fresh final authorization may each validate one cold store.
|
||||
expect(unrelatedParses).toBeLessThanOrEqual(unrelatedCount * 4);
|
||||
} finally {
|
||||
parseSpy.mockRestore();
|
||||
materializeSpy.mockRestore();
|
||||
await yielded;
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it.each([
|
||||
"unpublished revocation",
|
||||
"published grant",
|
||||
"unpublished grant",
|
||||
"unpublished creation",
|
||||
"registry restart grant",
|
||||
] as const)("rereads task access after a yielded %s", async (change) => {
|
||||
const config = rolePolicyConfig();
|
||||
const profileId = ensureProfileForEmail("task-viewer@example.test").id;
|
||||
const changingKey = "agent:main:changing-task-access";
|
||||
const stableKey = "agent:main:stable-task-access";
|
||||
const published = change === "published grant";
|
||||
const grant = change !== "unpublished revocation";
|
||||
const registryRestart = change === "registry restart grant";
|
||||
const changingIndex = !published && grant && !registryRestart ? 64 : 0;
|
||||
const entry = {
|
||||
sessionId: "changing-task-access",
|
||||
updatedAt: 1,
|
||||
visibility: grant ? ("draft" as const) : ("shared" as const),
|
||||
};
|
||||
if (change !== "unpublished creation") {
|
||||
await upsertSessionEntryCore({ agentId: "main", sessionKey: changingKey }, entry);
|
||||
}
|
||||
await upsertSessionEntryCore(
|
||||
{ agentId: "main", sessionKey: stableKey },
|
||||
{ sessionId: "stable-task-access", updatedAt: 1, visibility: "shared" },
|
||||
);
|
||||
const tasks = Array.from({ length: 65 }, (_, index) =>
|
||||
createSnapshotTask({
|
||||
taskId: `access-task-${index}`,
|
||||
requesterSessionKey: index === changingIndex ? changingKey : stableKey,
|
||||
requesterAgentId: "main",
|
||||
ownerKey: index === changingIndex ? changingKey : stableKey,
|
||||
lastEventAt: index === changingIndex ? 10_000 : 2_000 + index,
|
||||
}),
|
||||
);
|
||||
saveTaskRegistryStateToSqlite({
|
||||
tasks: new Map(tasks.map((task) => [task.taskId, task])),
|
||||
deliveryStates: new Map(),
|
||||
});
|
||||
reloadTaskRegistryFromStore();
|
||||
const context = {
|
||||
getRuntimeConfig: () => config,
|
||||
broadcast: () => {},
|
||||
getSessionEventSubscriberConnIds: () => new Set<string>(),
|
||||
};
|
||||
const accessRevision = readGatewayAccessRevision();
|
||||
// Exercise both an already-selected requester and one not yet visited when the scan yields.
|
||||
const mutation = new Promise<void>((resolve, reject) => {
|
||||
setImmediate(() => {
|
||||
void (async () => {
|
||||
if (published) {
|
||||
const { calls, respond } = captureRespond();
|
||||
await expectDefined(
|
||||
sessionSharingHandlers["session.visibility.set"],
|
||||
"session.visibility.set handler",
|
||||
)({
|
||||
params: { sessionKey: changingKey, agentId: "main", visibility: "shared" },
|
||||
client: identifiedClient(["operator.admin"], profileId),
|
||||
context,
|
||||
respond,
|
||||
} as never);
|
||||
expect(calls[0]?.[0]).toBe(true);
|
||||
expect(readGatewayAccessRevision()).toBeGreaterThan(accessRevision);
|
||||
} else {
|
||||
await upsertSessionEntryCore(
|
||||
{ agentId: "main", sessionKey: changingKey },
|
||||
{ ...entry, visibility: grant ? "shared" : "draft", updatedAt: 2 },
|
||||
);
|
||||
expect(readGatewayAccessRevision()).toBe(accessRevision);
|
||||
if (registryRestart) {
|
||||
expect(deleteTaskRecordById("access-task-63")).toBe(true);
|
||||
}
|
||||
}
|
||||
})().then(resolve, reject);
|
||||
});
|
||||
});
|
||||
const [{ calls, payload }] = await Promise.all([
|
||||
runTaskHandler(
|
||||
"tasks.list",
|
||||
{ limit: 1 },
|
||||
config,
|
||||
identifiedClient(["operator.read"], profileId),
|
||||
context as never,
|
||||
),
|
||||
mutation,
|
||||
]);
|
||||
expect(calls[0]?.[0]).toBe(true);
|
||||
expect(payload?.tasks?.map((task) => task.id)).toEqual([
|
||||
grant ? `access-task-${changingIndex}` : "access-task-64",
|
||||
]);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,94 @@
|
||||
import { expectDefined } from "@openclaw/normalization-core";
|
||||
import type { TaskRecord } from "../../tasks/task-registry.types.js";
|
||||
import { tasksHandlers } from "./tasks.js";
|
||||
import type { GatewayClient, RespondFn } from "./types.js";
|
||||
|
||||
type TaskResponsePayload = {
|
||||
tasks?: Array<Record<string, unknown>>;
|
||||
task?: Record<string, unknown>;
|
||||
found?: boolean;
|
||||
cancelled?: boolean;
|
||||
nextCursor?: string;
|
||||
results?: Array<{ taskId?: string; ok?: boolean; reason?: string }>;
|
||||
};
|
||||
|
||||
export function identifiedClient(
|
||||
scopes: string[],
|
||||
profileId = "viewer@example.com",
|
||||
): GatewayClient {
|
||||
return {
|
||||
connId: `conn-${profileId}-${scopes.join("-")}`,
|
||||
connect: {
|
||||
minProtocol: 1,
|
||||
maxProtocol: 1,
|
||||
client: { id: "openclaw-control-ui", version: "test", platform: "test", mode: "webchat" },
|
||||
role: "operator",
|
||||
scopes,
|
||||
},
|
||||
authenticatedUserId: "viewer@example.com",
|
||||
authenticatedUserProfile: {
|
||||
profileId,
|
||||
displayName: null,
|
||||
hasAvatar: false,
|
||||
updatedAt: 1,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export function captureRespond() {
|
||||
const calls: Parameters<RespondFn>[] = [];
|
||||
const respond: RespondFn = (...args) => {
|
||||
calls.push(args);
|
||||
};
|
||||
return { calls, respond };
|
||||
}
|
||||
|
||||
export function createContext(config: Record<string, unknown> = {}) {
|
||||
return {
|
||||
getRuntimeConfig: () => config,
|
||||
} as never;
|
||||
}
|
||||
|
||||
export function createSnapshotTask(overrides: Partial<TaskRecord>): TaskRecord {
|
||||
return {
|
||||
taskId: "task-snapshot",
|
||||
runtime: "cli",
|
||||
requesterSessionKey: "agent:main:main",
|
||||
ownerKey: "agent:main:main",
|
||||
scopeKind: "session",
|
||||
runId: "run-snapshot",
|
||||
task: "Snapshot task",
|
||||
status: "running",
|
||||
deliveryStatus: "pending",
|
||||
notifyPolicy: "done_only",
|
||||
createdAt: 1_000,
|
||||
startedAt: 1_010,
|
||||
lastEventAt: 1_010,
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
export async function runTaskHandler(
|
||||
method: "tasks.list" | "tasks.get" | "tasks.cancel" | "tasks.retry" | "tasks.dismiss",
|
||||
params: Record<string, unknown>,
|
||||
config: Record<string, unknown> = {},
|
||||
client: GatewayClient | null = null,
|
||||
context = createContext(config),
|
||||
) {
|
||||
const { calls, respond } = captureRespond();
|
||||
await expectDefined(
|
||||
tasksHandlers[method],
|
||||
"tasksHandlers[method] test invariant",
|
||||
)({
|
||||
req: { type: "req", id: `req-${method}`, method },
|
||||
params,
|
||||
respond,
|
||||
context,
|
||||
client,
|
||||
isWebchatConnect: () => false,
|
||||
});
|
||||
return {
|
||||
calls,
|
||||
payload: calls[0]?.[1] as TaskResponsePayload | undefined,
|
||||
};
|
||||
}
|
||||
@@ -34,8 +34,12 @@ import {
|
||||
setTaskRegistryControlRuntimeForTests,
|
||||
} from "../../tasks/task-runtime.test-helpers.js";
|
||||
import { captureEnv, setTestEnvValue } from "../../test-utils/env.js";
|
||||
import { tasksHandlers } from "./tasks.js";
|
||||
import type { GatewayClient, RespondFn } from "./types.js";
|
||||
import {
|
||||
createContext,
|
||||
createSnapshotTask,
|
||||
identifiedClient,
|
||||
runTaskHandler,
|
||||
} from "./tasks.test-helpers.js";
|
||||
|
||||
const stateDirEnvSnapshot = captureEnv(["OPENCLAW_STATE_DIR"]);
|
||||
const cancelSessionMock = vi.fn();
|
||||
@@ -44,14 +48,6 @@ const mainSessionTaskScope = {
|
||||
ownerKey: "agent:main:main",
|
||||
scopeKind: "session",
|
||||
} as const;
|
||||
type TaskResponsePayload = {
|
||||
tasks?: Array<Record<string, unknown>>;
|
||||
task?: Record<string, unknown>;
|
||||
found?: boolean;
|
||||
cancelled?: boolean;
|
||||
nextCursor?: string;
|
||||
results?: Array<{ taskId?: string; ok?: boolean; reason?: string }>;
|
||||
};
|
||||
|
||||
let stateDir: string;
|
||||
|
||||
@@ -88,84 +84,6 @@ afterEach(async () => {
|
||||
await fs.rm(stateDir, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
function identifiedClient(scopes: string[], profileId = "viewer@example.com"): GatewayClient {
|
||||
return {
|
||||
connId: `conn-${profileId}-${scopes.join("-")}`,
|
||||
connect: {
|
||||
minProtocol: 1,
|
||||
maxProtocol: 1,
|
||||
client: { id: "openclaw-control-ui", version: "test", platform: "test", mode: "webchat" },
|
||||
role: "operator",
|
||||
scopes,
|
||||
},
|
||||
authenticatedUserId: "viewer@example.com",
|
||||
authenticatedUserProfile: {
|
||||
profileId,
|
||||
displayName: null,
|
||||
hasAvatar: false,
|
||||
updatedAt: 1,
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function captureRespond() {
|
||||
const calls: Parameters<RespondFn>[] = [];
|
||||
const respond: RespondFn = (...args) => {
|
||||
calls.push(args);
|
||||
};
|
||||
return { calls, respond };
|
||||
}
|
||||
|
||||
function createContext(config: Record<string, unknown> = {}) {
|
||||
return {
|
||||
getRuntimeConfig: () => config,
|
||||
} as never;
|
||||
}
|
||||
|
||||
function createSnapshotTask(overrides: Partial<TaskRecord>): TaskRecord {
|
||||
return {
|
||||
taskId: "task-snapshot",
|
||||
runtime: "cli",
|
||||
requesterSessionKey: "agent:main:main",
|
||||
ownerKey: "agent:main:main",
|
||||
scopeKind: "session",
|
||||
runId: "run-snapshot",
|
||||
task: "Snapshot task",
|
||||
status: "running",
|
||||
deliveryStatus: "pending",
|
||||
notifyPolicy: "done_only",
|
||||
createdAt: 1_000,
|
||||
startedAt: 1_010,
|
||||
lastEventAt: 1_010,
|
||||
...overrides,
|
||||
};
|
||||
}
|
||||
|
||||
async function runTaskHandler(
|
||||
method: "tasks.list" | "tasks.get" | "tasks.cancel" | "tasks.retry" | "tasks.dismiss",
|
||||
params: Record<string, unknown>,
|
||||
config: Record<string, unknown> = {},
|
||||
client: GatewayClient | null = null,
|
||||
context = createContext(config),
|
||||
) {
|
||||
const { calls, respond } = captureRespond();
|
||||
await expectDefined(
|
||||
tasksHandlers[method],
|
||||
"tasksHandlers[method] test invariant",
|
||||
)({
|
||||
req: { type: "req", id: `req-${method}`, method },
|
||||
params,
|
||||
respond,
|
||||
context,
|
||||
client,
|
||||
isWebchatConnect: () => false,
|
||||
});
|
||||
return {
|
||||
calls,
|
||||
payload: calls[0]?.[1] as TaskResponsePayload | undefined,
|
||||
};
|
||||
}
|
||||
|
||||
async function getTaskPayload(taskId: string) {
|
||||
const { calls, payload } = await runTaskHandler("tasks.get", { taskId });
|
||||
expect(calls[0]?.[0]).toBe(true);
|
||||
|
||||
@@ -22,7 +22,10 @@ import { getTaskById, listTaskRecordPage } from "../../tasks/runtime-internal.js
|
||||
import type { TaskRecord, TaskStatus } from "../../tasks/task-registry.types.js";
|
||||
import { readGatewayAccessRevision } from "../gateway-access-revision.js";
|
||||
import { resolveRequestedSessionAgentId } from "../session-request-agent.js";
|
||||
import { canAccessTaskRequesterSession } from "../task-session-access.js";
|
||||
import {
|
||||
canAccessTaskRequesterSession,
|
||||
prepareTaskSessionReadFilter,
|
||||
} from "../task-session-access.js";
|
||||
import { mapTaskSummary } from "./task-summary.js";
|
||||
import type { GatewayRequestHandlers } from "./types.js";
|
||||
import { assertValidParams } from "./validation.js";
|
||||
@@ -194,8 +197,8 @@ export const tasksHandlers: GatewayRequestHandlers = {
|
||||
}
|
||||
// Selection stays inside the registry so ordering applies before pagination
|
||||
// and only the bounded wire page pays for defensive record cloning.
|
||||
const canReadTask = (task: Readonly<TaskRecord>) =>
|
||||
canAccessTaskRequesterSession({ cfg, client, task });
|
||||
const prepareFilter = (tasks: readonly Readonly<TaskRecord>[]) =>
|
||||
prepareTaskSessionReadFilter({ cfg, client }, tasks);
|
||||
const pageParams = {
|
||||
offset: cursor?.offset ?? 0,
|
||||
limit,
|
||||
@@ -205,7 +208,7 @@ export const tasksHandlers: GatewayRequestHandlers = {
|
||||
sessionKey,
|
||||
sessionAgentId,
|
||||
cfg,
|
||||
filter: canReadTask,
|
||||
prepareFilter,
|
||||
sortBy: params.sortBy,
|
||||
};
|
||||
// Page scans yield to active task updates. Restart the complete selection
|
||||
@@ -231,7 +234,7 @@ export const tasksHandlers: GatewayRequestHandlers = {
|
||||
// Recheck selected rows in the final synchronous response turn as well.
|
||||
if (
|
||||
accessRevision !== readGatewayAccessRevision() ||
|
||||
page.tasks.some((task) => !canReadTask(task))
|
||||
!page.tasks.every(prepareFilter(page.tasks))
|
||||
) {
|
||||
if (cursor) {
|
||||
invalidTaskListCursor(respond);
|
||||
|
||||
@@ -22,9 +22,10 @@ import {
|
||||
} from "./server-methods/gateway-client-identity.js";
|
||||
import type { GatewayClient } from "./server-methods/types.js";
|
||||
import { prepareSessionCreatorProfile } from "./session-creator.js";
|
||||
import type {
|
||||
GatewaySessionStoreCache,
|
||||
GatewaySessionStoreDiscoveryCache,
|
||||
import {
|
||||
resolveGatewaySessionStoreTargetsReadOnly,
|
||||
type GatewaySessionStoreCache,
|
||||
type GatewaySessionStoreDiscoveryCache,
|
||||
} from "./session-utils-store-lookup.js";
|
||||
import {
|
||||
resolveCanonicalSessionStoreMatchFromStoreKeys,
|
||||
@@ -89,6 +90,23 @@ export function resolveSessionSharingTarget(params: {
|
||||
...(params.storeCache ? { storeCache: params.storeCache } : {}),
|
||||
...(params.targetDiscoveryCache ? { targetDiscoveryCache: params.targetDiscoveryCache } : {}),
|
||||
});
|
||||
return toSessionSharingTarget(target);
|
||||
}
|
||||
|
||||
/** Fresh metadata for one synchronous batch; no authorization decisions are retained. */
|
||||
export function resolveSessionSharingTargets(params: {
|
||||
cfg: OpenClawConfig;
|
||||
targets: readonly { sessionKey: string; agentId?: string }[];
|
||||
}): Array<SessionSharingTarget | null> {
|
||||
return resolveGatewaySessionStoreTargetsReadOnly({
|
||||
cfg: params.cfg,
|
||||
targets: params.targets.map(({ sessionKey, agentId }) => ({ key: sessionKey, agentId })),
|
||||
}).map(toSessionSharingTarget);
|
||||
}
|
||||
|
||||
function toSessionSharingTarget(
|
||||
target: ReturnType<typeof resolveGatewaySessionStoreTargetWithStore>,
|
||||
): SessionSharingTarget | null {
|
||||
const match = resolveCanonicalSessionStoreMatchFromStoreKeys(target.store, target.storeKeys);
|
||||
return match
|
||||
? {
|
||||
|
||||
@@ -100,6 +100,7 @@ export {
|
||||
isSessionVisibilityAllowed,
|
||||
resolveSessionSharingRole,
|
||||
resolveSessionSharingTarget,
|
||||
resolveSessionSharingTargets,
|
||||
resolveSessionVisibility,
|
||||
} from "./session-sharing-policy.js";
|
||||
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import { expectDefined } from "@openclaw/normalization-core";
|
||||
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
|
||||
import { listAgentIds } from "../agents/agent-scope.js";
|
||||
import {
|
||||
@@ -10,9 +11,6 @@ import {
|
||||
} from "../config/sessions.js";
|
||||
import {
|
||||
listSessionChildEntriesReadOnly,
|
||||
listSessionEntriesCore as listAccessorSessionEntries,
|
||||
listSessionEntriesReadOnly as listAccessorSessionEntriesReadOnly,
|
||||
loadExactSessionEntryCandidates,
|
||||
type SessionEntryListScope,
|
||||
} from "../config/sessions/session-accessor.js";
|
||||
import { canonicalSessionKeyMigrationRequiredError } from "../config/sessions/session-canonical-key.js";
|
||||
@@ -33,6 +31,13 @@ import type {
|
||||
GatewaySessionStoreTarget,
|
||||
GatewaySessionStoreTargetWithStore,
|
||||
} from "./session-utils-contracts.js";
|
||||
import {
|
||||
loadGatewaySessionStoreReads,
|
||||
readGatewaySessionStore,
|
||||
type GatewaySessionStoreRead,
|
||||
type GatewaySessionStoreCache,
|
||||
} from "./session-utils-store-read.js";
|
||||
export type { GatewaySessionStoreCache } from "./session-utils-store-read.js";
|
||||
|
||||
function findCanonicalStoreMatch(
|
||||
store: Record<string, SessionEntry>,
|
||||
@@ -127,18 +132,6 @@ function resolveGatewaySessionStoreCandidates(
|
||||
return discovery;
|
||||
}
|
||||
|
||||
/**
|
||||
* Request-scoped store reuse.
|
||||
*
|
||||
* Sharing resolution runs once per listed row, and each run materialized every
|
||||
* entry of a candidate store, making `sessions.list` quadratic in entries. A
|
||||
* caller that resolves many keys against the same stores passes one cache so
|
||||
* each store is materialized once. Entries are shared across rows within that
|
||||
* request, so cached stores are read-only to their holder; the cache is never
|
||||
* process-global, so it cannot serve a later request stale rows.
|
||||
*/
|
||||
export type GatewaySessionStoreCache = Map<string, Record<string, SessionEntry>>;
|
||||
|
||||
/**
|
||||
* Sharing resolves every returned row, but store targets are stable within one request.
|
||||
* Keep discovery agent-scoped here or each row repeats registry probes and agent-root scans.
|
||||
@@ -180,91 +173,36 @@ export function createGatewaySessionStoreDiscoveryCache(params: {
|
||||
return cache;
|
||||
}
|
||||
|
||||
function loadGatewaySessionLookupStore(
|
||||
storePath: string,
|
||||
clone: boolean | undefined,
|
||||
agentId?: string,
|
||||
options: {
|
||||
readOnly?: boolean;
|
||||
cache?: GatewaySessionStoreCache;
|
||||
exactKeys?: readonly string[];
|
||||
projection?: SessionEntryListScope["projection"];
|
||||
} = {},
|
||||
): Record<string, SessionEntry> {
|
||||
const cache = options.cache;
|
||||
const cacheKey = cache
|
||||
? `${storePath}\u0000${agentId ?? ""}\u0000${clone === false ? "0" : "1"}\u0000${options.readOnly}\u0000${options.projection ?? "full"}\u0000${options.exactKeys?.join("\u0001") ?? ""}`
|
||||
: "";
|
||||
if (cache) {
|
||||
const cached = cache.get(cacheKey);
|
||||
if (cached) {
|
||||
return cached;
|
||||
}
|
||||
}
|
||||
const loaded = loadGatewaySessionLookupStoreUncached(storePath, clone, agentId, options);
|
||||
cache?.set(cacheKey, loaded);
|
||||
return loaded;
|
||||
}
|
||||
|
||||
function loadGatewaySessionLookupStoreUncached(
|
||||
storePath: string,
|
||||
clone: boolean | undefined,
|
||||
agentId?: string,
|
||||
options: {
|
||||
exactKeys?: readonly string[];
|
||||
readOnly?: boolean;
|
||||
projection?: SessionEntryListScope["projection"];
|
||||
} = {},
|
||||
): Record<string, SessionEntry> {
|
||||
try {
|
||||
if (options.exactKeys) {
|
||||
// Borrowed listing views and probes never create stores; ordinary owned reads may.
|
||||
return Object.fromEntries(
|
||||
loadExactSessionEntryCandidates({
|
||||
...(agentId ? { agentId } : {}),
|
||||
clone: false,
|
||||
projection: options.projection,
|
||||
sessionKeys: options.exactKeys,
|
||||
readOnly: options.readOnly !== false || clone === false,
|
||||
storePath,
|
||||
}).map(({ sessionKey, entry }) => [sessionKey, entry]),
|
||||
);
|
||||
}
|
||||
const listEntries = options.readOnly
|
||||
? listAccessorSessionEntriesReadOnly
|
||||
: listAccessorSessionEntries;
|
||||
return Object.fromEntries(
|
||||
listEntries({
|
||||
...(agentId ? { agentId } : {}),
|
||||
...(clone === false ? { clone: false } : {}),
|
||||
...(options.projection ? { projection: options.projection } : {}),
|
||||
storePath,
|
||||
}).map(({ sessionKey, entry }) => [sessionKey, entry]),
|
||||
);
|
||||
} catch {
|
||||
return {};
|
||||
}
|
||||
}
|
||||
|
||||
function resolveGatewaySessionStoreLookup(params: {
|
||||
type GatewaySessionStoreLookupParams = {
|
||||
cfg: OpenClawConfig;
|
||||
key: string;
|
||||
canonicalKey: string;
|
||||
agentId: string;
|
||||
agentId?: string;
|
||||
clone?: boolean;
|
||||
initialStore?: Record<string, SessionEntry>;
|
||||
projection?: SessionEntryListScope["projection"];
|
||||
readOnly?: boolean;
|
||||
exactRead?: boolean;
|
||||
deferCanonicalValidation?: boolean;
|
||||
includeStoreChildEntries?: boolean;
|
||||
store?: Record<string, SessionEntry>;
|
||||
storeCache?: GatewaySessionStoreCache;
|
||||
targetDiscoveryCache?: GatewaySessionStoreDiscoveryCache;
|
||||
}): {
|
||||
};
|
||||
|
||||
type GatewaySessionStorePlan<T> = {
|
||||
reads: GatewaySessionStoreRead[];
|
||||
resolve: () => T;
|
||||
};
|
||||
|
||||
type GatewaySessionStoreLookup = {
|
||||
storePath: string;
|
||||
store: Record<string, SessionEntry>;
|
||||
match: { entry: SessionEntry; key: string } | undefined;
|
||||
canonicalValidationError?: Error;
|
||||
} {
|
||||
};
|
||||
|
||||
function prepareGatewaySessionStoreLookup(
|
||||
params: GatewaySessionStoreLookupParams & { canonicalKey: string; agentId: string },
|
||||
): GatewaySessionStorePlan<GatewaySessionStoreLookup> {
|
||||
const scanTargets = buildGatewaySessionStoreScanTargets(params);
|
||||
const discovery = resolveGatewaySessionStoreCandidates(
|
||||
params.cfg,
|
||||
@@ -279,67 +217,66 @@ function resolveGatewaySessionStoreLookup(params: {
|
||||
? [fallback, ...existing.filter((target) => target.storePath !== fallback.storePath)]
|
||||
: existing;
|
||||
if (candidates.length === 0) {
|
||||
// Discovery is read-only. Only configured agents may cross the fallback edge that creates a
|
||||
// missing SQLite store; retired/manual agents must already have a discovered store.
|
||||
// Retired/manual agents require an existing discovered store; lookup never creates one.
|
||||
return {
|
||||
storePath: fallback.storePath,
|
||||
store: {},
|
||||
match: undefined,
|
||||
reads: [],
|
||||
resolve: () => ({ storePath: fallback.storePath, store: {}, match: undefined }),
|
||||
};
|
||||
}
|
||||
const loadStore = (target: SessionStoreTarget) =>
|
||||
loadGatewaySessionLookupStore(target.storePath, params.clone, target.agentId, {
|
||||
const reads = candidates.map((target, index): GatewaySessionStoreRead => ({
|
||||
storePath: target.storePath,
|
||||
agentId: target.agentId,
|
||||
clone: params.clone,
|
||||
options: {
|
||||
readOnly: configured ? params.readOnly : true,
|
||||
...(params.exactRead ? { exactKeys: scanTargets } : {}),
|
||||
...(params.projection ? { projection: params.projection } : {}),
|
||||
...(params.storeCache ? { cache: params.storeCache } : {}),
|
||||
});
|
||||
const firstCandidate = candidates[0] ?? fallback;
|
||||
let selectedStorePath = firstCandidate.storePath;
|
||||
let selectedStore =
|
||||
params.initialStore && firstCandidate.storePath === fallback.storePath
|
||||
? params.initialStore
|
||||
: loadStore(firstCandidate);
|
||||
let canonicalValidationError: Error | undefined;
|
||||
const recordCanonicalError = params.deferCanonicalValidation
|
||||
? (error: Error) => {
|
||||
canonicalValidationError ??= error;
|
||||
}
|
||||
: undefined;
|
||||
let selectedMatch = findCanonicalStoreMatch(selectedStore, scanTargets, recordCanonicalError);
|
||||
|
||||
for (let index = 1; index < candidates.length; index += 1) {
|
||||
const candidate = candidates[index];
|
||||
if (!candidate) {
|
||||
continue;
|
||||
}
|
||||
const store = loadStore(candidate);
|
||||
const match = findCanonicalStoreMatch(store, scanTargets, recordCanonicalError);
|
||||
if (!match) {
|
||||
continue;
|
||||
}
|
||||
if (selectedMatch) {
|
||||
const error = canonicalSessionKeyMigrationRequiredError(
|
||||
`duplicate rows resolve to canonical session key ${params.canonicalKey}`,
|
||||
);
|
||||
if (!recordCanonicalError) {
|
||||
throw error;
|
||||
}
|
||||
recordCanonicalError(error);
|
||||
if (match.key !== params.canonicalKey || selectedMatch.key === params.canonicalKey) {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
selectedStorePath = candidate.storePath;
|
||||
selectedStore = store;
|
||||
selectedMatch = match;
|
||||
}
|
||||
|
||||
},
|
||||
store: index === 0 && target.storePath === fallback.storePath ? params.store : undefined,
|
||||
}));
|
||||
return {
|
||||
storePath: selectedStorePath,
|
||||
store: selectedStore,
|
||||
match: selectedMatch,
|
||||
...(canonicalValidationError ? { canonicalValidationError } : {}),
|
||||
reads,
|
||||
resolve: () => {
|
||||
const first = expectDefined(reads[0], "first configured or discovered session store");
|
||||
let selectedStorePath = first.storePath;
|
||||
let selectedStore = readGatewaySessionStore(first);
|
||||
let canonicalValidationError: Error | undefined;
|
||||
const recordCanonicalError = params.deferCanonicalValidation
|
||||
? (error: Error) => {
|
||||
canonicalValidationError ??= error;
|
||||
}
|
||||
: undefined;
|
||||
let selectedMatch = findCanonicalStoreMatch(selectedStore, scanTargets, recordCanonicalError);
|
||||
for (const candidate of reads.slice(1)) {
|
||||
const store = readGatewaySessionStore(candidate);
|
||||
const match = findCanonicalStoreMatch(store, scanTargets, recordCanonicalError);
|
||||
if (!match) {
|
||||
continue;
|
||||
}
|
||||
if (selectedMatch) {
|
||||
const error = canonicalSessionKeyMigrationRequiredError(
|
||||
`duplicate rows resolve to canonical session key ${params.canonicalKey}`,
|
||||
);
|
||||
if (!recordCanonicalError) {
|
||||
throw error;
|
||||
}
|
||||
recordCanonicalError(error);
|
||||
if (match.key !== params.canonicalKey || selectedMatch.key === params.canonicalKey) {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
selectedStorePath = candidate.storePath;
|
||||
selectedStore = store;
|
||||
selectedMatch = match;
|
||||
}
|
||||
return {
|
||||
storePath: selectedStorePath,
|
||||
store: selectedStore,
|
||||
match: selectedMatch,
|
||||
...(canonicalValidationError ? { canonicalValidationError } : {}),
|
||||
};
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
@@ -347,17 +284,9 @@ function isAgentScopedSentinelSessionKey(canonicalKey: string): boolean {
|
||||
return canonicalKey === "global" || canonicalKey === "unknown";
|
||||
}
|
||||
|
||||
function resolveExplicitDeletedLegacyMainStoreTarget(params: {
|
||||
cfg: OpenClawConfig;
|
||||
key: string;
|
||||
clone?: boolean;
|
||||
deferCanonicalValidation?: boolean;
|
||||
projection?: SessionEntryListScope["projection"];
|
||||
readOnly?: boolean;
|
||||
exactRead?: boolean;
|
||||
storeCache?: GatewaySessionStoreCache;
|
||||
targetDiscoveryCache?: GatewaySessionStoreDiscoveryCache;
|
||||
}): GatewaySessionStoreTargetWithStore | null {
|
||||
function prepareExplicitDeletedLegacyMainStoreTarget(
|
||||
params: GatewaySessionStoreLookupParams,
|
||||
): GatewaySessionStorePlan<GatewaySessionStoreTargetWithStore | null> | null {
|
||||
const parsed = parseAgentSessionKey(params.key);
|
||||
const legacyAgentId = normalizeAgentId(parsed?.agentId);
|
||||
if (
|
||||
@@ -368,120 +297,96 @@ function resolveExplicitDeletedLegacyMainStoreTarget(params: {
|
||||
) {
|
||||
return null;
|
||||
}
|
||||
|
||||
// Only preserve agent:main:* when it is backed by a discovered deleted-main store.
|
||||
// Shared-store legacy aliases should continue remapping to the configured default agent.
|
||||
// Deleted-main discovery precedes normal aliases; only a real matching row keeps this owner.
|
||||
const canonicalKey = resolveStoredSessionKeyForAgentStore({
|
||||
cfg: params.cfg,
|
||||
agentId: legacyAgentId,
|
||||
sessionKey: params.key,
|
||||
});
|
||||
const agentMainKey = resolveAgentMainSessionKey({ cfg: params.cfg, agentId: legacyAgentId });
|
||||
const legacyAgentMainKey = `agent:${legacyAgentId}:main`;
|
||||
const lookupSeeds = Array.from(
|
||||
new Set([params.key, canonicalKey, agentMainKey, legacyAgentMainKey]),
|
||||
new Set([params.key, canonicalKey, agentMainKey, `agent:${legacyAgentId}:main`]),
|
||||
);
|
||||
let best:
|
||||
| {
|
||||
storePath: string;
|
||||
store: Record<string, SessionEntry>;
|
||||
match: { entry: SessionEntry; key: string };
|
||||
}
|
||||
| undefined;
|
||||
const { existing } = resolveGatewaySessionStoreCandidates(
|
||||
params.cfg,
|
||||
legacyAgentId,
|
||||
params.targetDiscoveryCache,
|
||||
);
|
||||
let canonicalValidationError: Error | undefined;
|
||||
const recordCanonicalError = params.deferCanonicalValidation
|
||||
? (error: Error) => {
|
||||
canonicalValidationError ??= error;
|
||||
}
|
||||
: undefined;
|
||||
for (const target of existing) {
|
||||
if (target.agentId !== legacyAgentId) {
|
||||
continue;
|
||||
}
|
||||
const store = loadGatewaySessionLookupStore(target.storePath, params.clone, target.agentId, {
|
||||
readOnly: true,
|
||||
...(params.exactRead ? { exactKeys: lookupSeeds } : {}),
|
||||
...(params.projection ? { projection: params.projection } : {}),
|
||||
...(params.storeCache ? { cache: params.storeCache } : {}),
|
||||
});
|
||||
const match = findCanonicalStoreMatch(store, lookupSeeds, recordCanonicalError);
|
||||
if (!match) {
|
||||
continue;
|
||||
}
|
||||
if (best) {
|
||||
const error = canonicalSessionKeyMigrationRequiredError(
|
||||
`duplicate rows resolve to canonical session key ${canonicalKey}`,
|
||||
);
|
||||
if (!recordCanonicalError) {
|
||||
throw error;
|
||||
}
|
||||
recordCanonicalError(error);
|
||||
}
|
||||
if (!best || (match.entry.updatedAt ?? 0) >= (best.match.entry.updatedAt ?? 0)) {
|
||||
best = { storePath: target.storePath, store, match };
|
||||
}
|
||||
}
|
||||
if (!best) {
|
||||
return null;
|
||||
}
|
||||
|
||||
const storeKeys = new Set<string>([canonicalKey]);
|
||||
if (params.key !== canonicalKey) {
|
||||
storeKeys.add(params.key);
|
||||
}
|
||||
storeKeys.add(best.match.key);
|
||||
for (const seed of lookupSeeds) {
|
||||
storeKeys.add(seed);
|
||||
}
|
||||
const reads = existing
|
||||
.filter((target) => target.agentId === legacyAgentId)
|
||||
.map((target): GatewaySessionStoreRead => ({
|
||||
storePath: target.storePath,
|
||||
clone: params.clone,
|
||||
agentId: target.agentId,
|
||||
options: {
|
||||
readOnly: true,
|
||||
...(params.exactRead ? { exactKeys: lookupSeeds } : {}),
|
||||
...(params.projection ? { projection: params.projection } : {}),
|
||||
...(params.storeCache ? { cache: params.storeCache } : {}),
|
||||
},
|
||||
}));
|
||||
return {
|
||||
agentId: legacyAgentId,
|
||||
storePath: best.storePath,
|
||||
canonicalKey,
|
||||
storeKeys: Array.from(storeKeys),
|
||||
store: best.store,
|
||||
...(canonicalValidationError ? { canonicalValidationError } : {}),
|
||||
reads,
|
||||
resolve: () => {
|
||||
let best:
|
||||
| {
|
||||
storePath: string;
|
||||
store: Record<string, SessionEntry>;
|
||||
match: { entry: SessionEntry; key: string };
|
||||
}
|
||||
| undefined;
|
||||
let canonicalValidationError: Error | undefined;
|
||||
const recordCanonicalError = params.deferCanonicalValidation
|
||||
? (error: Error) => {
|
||||
canonicalValidationError ??= error;
|
||||
}
|
||||
: undefined;
|
||||
for (const target of reads) {
|
||||
const store = readGatewaySessionStore(target);
|
||||
const match = findCanonicalStoreMatch(store, lookupSeeds, recordCanonicalError);
|
||||
if (!match) {
|
||||
continue;
|
||||
}
|
||||
if (best) {
|
||||
const error = canonicalSessionKeyMigrationRequiredError(
|
||||
`duplicate rows resolve to canonical session key ${canonicalKey}`,
|
||||
);
|
||||
if (!recordCanonicalError) {
|
||||
throw error;
|
||||
}
|
||||
recordCanonicalError(error);
|
||||
}
|
||||
if (!best || (match.entry.updatedAt ?? 0) >= (best.match.entry.updatedAt ?? 0)) {
|
||||
best = { storePath: target.storePath, store, match };
|
||||
}
|
||||
}
|
||||
if (!best) {
|
||||
return null;
|
||||
}
|
||||
const storeKeys = new Set<string>([canonicalKey]);
|
||||
if (params.key !== canonicalKey) {
|
||||
storeKeys.add(params.key);
|
||||
}
|
||||
storeKeys.add(best.match.key);
|
||||
for (const seed of lookupSeeds) {
|
||||
storeKeys.add(seed);
|
||||
}
|
||||
return {
|
||||
agentId: legacyAgentId,
|
||||
storePath: best.storePath,
|
||||
canonicalKey,
|
||||
storeKeys: Array.from(storeKeys),
|
||||
store: best.store,
|
||||
...(canonicalValidationError ? { canonicalValidationError } : {}),
|
||||
};
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export function resolveGatewaySessionStoreTargetWithStore(params: {
|
||||
cfg: OpenClawConfig;
|
||||
key: string;
|
||||
agentId?: string;
|
||||
clone?: boolean;
|
||||
projection?: SessionEntryListScope["projection"];
|
||||
readOnly?: boolean;
|
||||
exactRead?: boolean;
|
||||
deferCanonicalValidation?: boolean;
|
||||
includeStoreChildEntries?: boolean;
|
||||
store?: Record<string, SessionEntry>;
|
||||
storeCache?: GatewaySessionStoreCache;
|
||||
targetDiscoveryCache?: GatewaySessionStoreDiscoveryCache;
|
||||
}): GatewaySessionStoreTargetWithStore {
|
||||
const key = normalizeOptionalString(params.key) ?? "";
|
||||
const explicitDeletedMainTarget = resolveExplicitDeletedLegacyMainStoreTarget({
|
||||
cfg: params.cfg,
|
||||
key,
|
||||
clone: params.clone,
|
||||
...(params.deferCanonicalValidation ? { deferCanonicalValidation: true } : {}),
|
||||
readOnly: params.readOnly,
|
||||
exactRead: params.exactRead,
|
||||
...(params.projection ? { projection: params.projection } : {}),
|
||||
...(params.storeCache ? { storeCache: params.storeCache } : {}),
|
||||
...(params.targetDiscoveryCache ? { targetDiscoveryCache: params.targetDiscoveryCache } : {}),
|
||||
});
|
||||
if (explicitDeletedMainTarget) {
|
||||
return includeDirectChildEntries(
|
||||
explicitDeletedMainTarget,
|
||||
params.includeStoreChildEntries,
|
||||
params.projection,
|
||||
);
|
||||
}
|
||||
|
||||
function prepareGatewaySessionStoreTarget(
|
||||
params: GatewaySessionStoreLookupParams,
|
||||
): GatewaySessionStorePlan<GatewaySessionStoreTargetWithStore> {
|
||||
const key = params.key;
|
||||
const requestedAgentId = normalizeOptionalString(params.agentId);
|
||||
const canonicalKey = resolveSessionStoreKey({
|
||||
cfg: params.cfg,
|
||||
@@ -495,73 +400,100 @@ export function resolveGatewaySessionStoreTargetWithStore(params: {
|
||||
: resolveSessionStoreAgentId(params.cfg, canonicalKey);
|
||||
if (isIncognitoSessionKey(canonicalKey)) {
|
||||
const storePath = resolveIncognitoOpenClawAgentSqlitePath({ agentId });
|
||||
// Session resolution may receive arbitrary stale keys; only creation/write
|
||||
// owners may materialize the process-lifetime incognito database.
|
||||
const store = loadGatewaySessionLookupStore(storePath, params.clone, agentId, {
|
||||
readOnly: true,
|
||||
...(params.exactRead ? { exactKeys: [canonicalKey] } : {}),
|
||||
...(params.projection ? { projection: params.projection } : {}),
|
||||
...(params.storeCache ? { cache: params.storeCache } : {}),
|
||||
});
|
||||
return includeDirectChildEntries(
|
||||
{
|
||||
const read: GatewaySessionStoreRead = {
|
||||
storePath,
|
||||
agentId,
|
||||
clone: params.clone,
|
||||
options: {
|
||||
// Arbitrary stale keys must not materialize process-lifetime incognito state.
|
||||
readOnly: true,
|
||||
...(params.exactRead ? { exactKeys: [canonicalKey] } : {}),
|
||||
...(params.projection ? { projection: params.projection } : {}),
|
||||
...(params.storeCache ? { cache: params.storeCache } : {}),
|
||||
},
|
||||
};
|
||||
return {
|
||||
reads: [read],
|
||||
resolve: () => ({
|
||||
agentId,
|
||||
storePath,
|
||||
canonicalKey,
|
||||
storeKeys: [canonicalKey],
|
||||
store,
|
||||
},
|
||||
params.includeStoreChildEntries,
|
||||
params.projection,
|
||||
);
|
||||
store: readGatewaySessionStore(read),
|
||||
}),
|
||||
};
|
||||
}
|
||||
const { canonicalValidationError, storePath, store } = resolveGatewaySessionStoreLookup({
|
||||
cfg: params.cfg,
|
||||
key,
|
||||
canonicalKey,
|
||||
agentId,
|
||||
clone: params.clone,
|
||||
readOnly: params.readOnly,
|
||||
exactRead: params.exactRead,
|
||||
deferCanonicalValidation: params.deferCanonicalValidation,
|
||||
initialStore: params.store,
|
||||
...(params.projection ? { projection: params.projection } : {}),
|
||||
...(params.storeCache ? { storeCache: params.storeCache } : {}),
|
||||
...(params.targetDiscoveryCache ? { targetDiscoveryCache: params.targetDiscoveryCache } : {}),
|
||||
});
|
||||
if (canonicalKey === "global" || canonicalKey === "unknown") {
|
||||
const storeKeys = key && key !== canonicalKey ? [canonicalKey, key] : [key];
|
||||
return includeDirectChildEntries(
|
||||
{
|
||||
const lookup = prepareGatewaySessionStoreLookup({ ...params, canonicalKey, agentId });
|
||||
return {
|
||||
reads: lookup.reads,
|
||||
resolve: () => {
|
||||
const { canonicalValidationError, storePath, store } = lookup.resolve();
|
||||
const storeKeys = isAgentScopedSentinelSessionKey(canonicalKey)
|
||||
? key && key !== canonicalKey
|
||||
? [canonicalKey, key]
|
||||
: [key]
|
||||
: Array.from(
|
||||
new Set(
|
||||
buildGatewaySessionStoreScanTargets({ cfg: params.cfg, key, canonicalKey, agentId }),
|
||||
),
|
||||
);
|
||||
return {
|
||||
agentId,
|
||||
storePath,
|
||||
canonicalKey,
|
||||
storeKeys,
|
||||
store,
|
||||
...(canonicalValidationError ? { canonicalValidationError } : {}),
|
||||
},
|
||||
params.includeStoreChildEntries,
|
||||
params.projection,
|
||||
);
|
||||
}
|
||||
|
||||
const storeKeys = new Set<string>(
|
||||
buildGatewaySessionStoreScanTargets({ cfg: params.cfg, key, canonicalKey, agentId }),
|
||||
);
|
||||
return includeDirectChildEntries(
|
||||
{
|
||||
agentId,
|
||||
storePath,
|
||||
canonicalKey,
|
||||
storeKeys: Array.from(storeKeys),
|
||||
store,
|
||||
...(canonicalValidationError ? { canonicalValidationError } : {}),
|
||||
};
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export function resolveGatewaySessionStoreTargetWithStore(
|
||||
params: GatewaySessionStoreLookupParams,
|
||||
): GatewaySessionStoreTargetWithStore {
|
||||
const normalized = { ...params, key: normalizeOptionalString(params.key) ?? "" };
|
||||
const deletedMain = prepareExplicitDeletedLegacyMainStoreTarget(normalized)?.resolve();
|
||||
return includeDirectChildEntries(
|
||||
deletedMain ?? prepareGatewaySessionStoreTarget(normalized).resolve(),
|
||||
params.includeStoreChildEntries,
|
||||
params.projection,
|
||||
);
|
||||
}
|
||||
|
||||
/** Resolve one synchronous set of logical metadata targets using exact grouped reads. */
|
||||
export function resolveGatewaySessionStoreTargetsReadOnly(params: {
|
||||
cfg: OpenClawConfig;
|
||||
targets: readonly { key: string; agentId?: string }[];
|
||||
}): GatewaySessionStoreTargetWithStore[] {
|
||||
const targetDiscoveryCache: GatewaySessionStoreDiscoveryCache = new Map();
|
||||
const requests = params.targets.map((target) => {
|
||||
const lookup: GatewaySessionStoreLookupParams = {
|
||||
...target,
|
||||
key: normalizeOptionalString(target.key) ?? "",
|
||||
cfg: params.cfg,
|
||||
clone: false,
|
||||
readOnly: true,
|
||||
exactRead: true,
|
||||
projection: "list",
|
||||
targetDiscoveryCache,
|
||||
};
|
||||
return { lookup, legacy: prepareExplicitDeletedLegacyMainStoreTarget(lookup) };
|
||||
});
|
||||
loadGatewaySessionStoreReads(requests.flatMap((request) => request.legacy?.reads ?? []));
|
||||
// A successful deleted-main lookup must not open or validate a normal fallback store.
|
||||
const selected = requests.map((request) => {
|
||||
const target = request.legacy?.resolve();
|
||||
return target ? { target } : { plan: prepareGatewaySessionStoreTarget(request.lookup) };
|
||||
});
|
||||
loadGatewaySessionStoreReads(selected.flatMap((selection) => selection.plan?.reads ?? []));
|
||||
return selected.map(
|
||||
(selection) =>
|
||||
selection.target ??
|
||||
expectDefined(selection.plan, "unresolved logical session plan").resolve(),
|
||||
);
|
||||
}
|
||||
|
||||
function includeDirectChildEntries(
|
||||
target: GatewaySessionStoreTargetWithStore,
|
||||
include: boolean | undefined,
|
||||
|
||||
@@ -0,0 +1,128 @@
|
||||
import { expectDefined } from "@openclaw/normalization-core";
|
||||
import type { SessionEntry } from "../config/sessions.js";
|
||||
import {
|
||||
listSessionEntriesCore as listAccessorSessionEntries,
|
||||
listSessionEntriesReadOnly as listAccessorSessionEntriesReadOnly,
|
||||
loadExactSessionEntryCandidates,
|
||||
loadExactSessionEntryCandidatesReadOnlyBatch,
|
||||
type SessionEntryListScope,
|
||||
} from "../config/sessions/session-accessor.js";
|
||||
|
||||
/**
|
||||
* Request-scoped store reuse.
|
||||
*
|
||||
* Sharing resolution runs once per listed row, and each run materialized every
|
||||
* entry of a candidate store, making `sessions.list` quadratic in entries. A
|
||||
* caller that resolves many keys against the same stores passes one cache so
|
||||
* each store is materialized once. Entries are shared across rows within that
|
||||
* request, so cached stores are read-only to their holder; the cache is never
|
||||
* process-global, so it cannot serve a later request stale rows.
|
||||
*/
|
||||
export type GatewaySessionStoreCache = Map<string, Record<string, SessionEntry>>;
|
||||
|
||||
export type GatewaySessionStoreRead = {
|
||||
storePath: string;
|
||||
clone?: boolean;
|
||||
agentId?: string;
|
||||
options: NonNullable<Parameters<typeof loadGatewaySessionLookupStore>[3]>;
|
||||
store?: Record<string, SessionEntry>;
|
||||
};
|
||||
|
||||
/** Single-target resolution keeps its original lazy read and failure order. */
|
||||
export function readGatewaySessionStore(
|
||||
read: GatewaySessionStoreRead,
|
||||
): Record<string, SessionEntry> {
|
||||
return (read.store ??= loadGatewaySessionLookupStore(
|
||||
read.storePath,
|
||||
read.clone,
|
||||
read.agentId,
|
||||
read.options,
|
||||
));
|
||||
}
|
||||
|
||||
/** Populate exact logical lookups without materializing unrelated store entries. */
|
||||
export function loadGatewaySessionStoreReads(reads: readonly GatewaySessionStoreRead[]): void {
|
||||
const pending = reads.filter((read) => read.store === undefined);
|
||||
const results = loadExactSessionEntryCandidatesReadOnlyBatch(
|
||||
pending.map((read) => ({
|
||||
agentId: read.agentId,
|
||||
storePath: read.storePath,
|
||||
projection: read.options.projection,
|
||||
clone: false,
|
||||
sessionKeys: expectDefined(read.options.exactKeys, "exact batch lookup keys"),
|
||||
})),
|
||||
);
|
||||
for (const [index, read] of pending.entries()) {
|
||||
const result = expectDefined(results[index], "exact batch lookup result");
|
||||
// Preserve the existing per-logical-target unreadable-store behavior.
|
||||
read.store = result.ok
|
||||
? Object.fromEntries(result.value.map(({ sessionKey, entry }) => [sessionKey, entry]))
|
||||
: {};
|
||||
}
|
||||
}
|
||||
|
||||
function loadGatewaySessionLookupStore(
|
||||
storePath: string,
|
||||
clone: boolean | undefined,
|
||||
agentId?: string,
|
||||
options: {
|
||||
readOnly?: boolean;
|
||||
cache?: GatewaySessionStoreCache;
|
||||
exactKeys?: readonly string[];
|
||||
projection?: SessionEntryListScope["projection"];
|
||||
} = {},
|
||||
): Record<string, SessionEntry> {
|
||||
const cache = options.cache;
|
||||
const cacheKey = cache
|
||||
? `${storePath}\u0000${agentId ?? ""}\u0000${clone === false ? "0" : "1"}\u0000${options.readOnly}\u0000${options.projection ?? "full"}\u0000${options.exactKeys?.join("\u0001") ?? ""}`
|
||||
: "";
|
||||
if (cache) {
|
||||
const cached = cache.get(cacheKey);
|
||||
if (cached) {
|
||||
return cached;
|
||||
}
|
||||
}
|
||||
const loaded = loadGatewaySessionLookupStoreUncached(storePath, clone, agentId, options);
|
||||
cache?.set(cacheKey, loaded);
|
||||
return loaded;
|
||||
}
|
||||
|
||||
function loadGatewaySessionLookupStoreUncached(
|
||||
storePath: string,
|
||||
clone: boolean | undefined,
|
||||
agentId?: string,
|
||||
options: {
|
||||
exactKeys?: readonly string[];
|
||||
readOnly?: boolean;
|
||||
projection?: SessionEntryListScope["projection"];
|
||||
} = {},
|
||||
): Record<string, SessionEntry> {
|
||||
try {
|
||||
if (options.exactKeys) {
|
||||
// Borrowed listing views and probes never create stores; ordinary owned reads may.
|
||||
return Object.fromEntries(
|
||||
loadExactSessionEntryCandidates({
|
||||
...(agentId ? { agentId } : {}),
|
||||
clone: false,
|
||||
projection: options.projection,
|
||||
sessionKeys: options.exactKeys,
|
||||
readOnly: options.readOnly !== false || clone === false,
|
||||
storePath,
|
||||
}).map(({ sessionKey, entry }) => [sessionKey, entry]),
|
||||
);
|
||||
}
|
||||
const listEntries = options.readOnly
|
||||
? listAccessorSessionEntriesReadOnly
|
||||
: listAccessorSessionEntries;
|
||||
return Object.fromEntries(
|
||||
listEntries({
|
||||
...(agentId ? { agentId } : {}),
|
||||
...(clone === false ? { clone: false } : {}),
|
||||
...(options.projection ? { projection: options.projection } : {}),
|
||||
storePath,
|
||||
}).map(({ sessionKey, entry }) => [sessionKey, entry]),
|
||||
);
|
||||
} catch {
|
||||
return {};
|
||||
}
|
||||
}
|
||||
@@ -52,6 +52,7 @@ import {
|
||||
type GatewaySessionStoreDiscoveryCache,
|
||||
resolveGatewaySessionStoreTarget,
|
||||
resolveGatewaySessionStoreTargetWithStore,
|
||||
resolveGatewaySessionStoreTargetsReadOnly,
|
||||
} from "./session-utils-store-lookup.js";
|
||||
import {
|
||||
listAgentsForGateway,
|
||||
@@ -3236,6 +3237,14 @@ describe("gateway session utils", () => {
|
||||
|
||||
expect(target.storePath).toBe(path.resolve(fixedStorePath));
|
||||
expect(target.store["agent:ops:main"]?.sessionId).toBe("sess-fixed");
|
||||
expect(
|
||||
resolveGatewaySessionStoreTargetsReadOnly({ cfg, targets: [{ key: "agent:ops:main" }] }),
|
||||
).toMatchObject([
|
||||
{
|
||||
storePath: path.resolve(fixedStorePath),
|
||||
store: { "agent:ops:main": { sessionId: "sess-fixed" } },
|
||||
},
|
||||
]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -3293,6 +3302,49 @@ describe("gateway session utils", () => {
|
||||
});
|
||||
});
|
||||
|
||||
test("batched session targets preserve explicit sentinel owners and reject discovered collisions", async () => {
|
||||
await withStateDirEnv("session-utils-batch-owners-", async ({ stateDir }) => {
|
||||
const cfg = {
|
||||
session: { store: path.join(stateDir, "agents", "{agentId}", "sessions", "sessions.json") },
|
||||
agents: { ownership: "explicit", entries: { ops: {}, research: {} } },
|
||||
} satisfies OpenClawConfig;
|
||||
for (const agentId of ["ops", "research"]) {
|
||||
await replaceSessionEntry(
|
||||
{
|
||||
agentId,
|
||||
sessionKey: "global",
|
||||
storePath: cfg.session.store.replaceAll("{agentId}", agentId),
|
||||
},
|
||||
{ sessionId: `global-${agentId}`, updatedAt: 1 },
|
||||
);
|
||||
}
|
||||
expect(
|
||||
resolveGatewaySessionStoreTargetsReadOnly({
|
||||
cfg,
|
||||
targets: ["research", "ops"].map((agentId) => ({ key: "global", agentId })),
|
||||
}),
|
||||
).toMatchObject([
|
||||
{ agentId: "research", store: { global: { sessionId: "global-research" } } },
|
||||
{ agentId: "ops", store: { global: { sessionId: "global-ops" } } },
|
||||
]);
|
||||
for (const directory of ["Retired Agent", "retired-agent"]) {
|
||||
await seedSessionEntries(
|
||||
path.join(stateDir, "agents", directory, "sessions", "sessions.json"),
|
||||
{
|
||||
"agent:retired-agent:main": { sessionId: directory, updatedAt: 1 },
|
||||
},
|
||||
);
|
||||
}
|
||||
const key = "agent:retired-agent:main";
|
||||
expect(() =>
|
||||
resolveGatewaySessionStoreTargetWithStore({ cfg, key, readOnly: true, exactRead: true }),
|
||||
).toThrow("openclaw doctor --fix");
|
||||
expect(() => resolveGatewaySessionStoreTargetsReadOnly({ cfg, targets: [{ key }] })).toThrow(
|
||||
"openclaw doctor --fix",
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
test("resolveGatewaySessionStoreTarget finds a retired agent's row under another configured agent's template root", async () => {
|
||||
await withStateDirEnv("session-utils-retired-cross-root-", async ({ tempRoot }) => {
|
||||
const storesRoot = path.join(tempRoot, "stores");
|
||||
@@ -3329,6 +3381,14 @@ describe("gateway session utils", () => {
|
||||
|
||||
expect(target.storePath).toBe(path.resolve(retiredStorePath));
|
||||
expect(target.store["agent:old:main"]?.sessionId).toBe("sess-retired-cross-root");
|
||||
expect(
|
||||
resolveGatewaySessionStoreTargetsReadOnly({ cfg, targets: [{ key: "agent:old:main" }] }),
|
||||
).toMatchObject([
|
||||
{
|
||||
storePath: path.resolve(retiredStorePath),
|
||||
store: { "agent:old:main": { sessionId: "sess-retired-cross-root" } },
|
||||
},
|
||||
]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -3365,6 +3425,13 @@ describe("gateway session utils", () => {
|
||||
expect(sqlitePath).toBeDefined();
|
||||
expect(fs.existsSync(sqlitePath!)).toBe(false);
|
||||
expect(fs.readdirSync(retiredSessionsDir)).toEqual(["sessions.json"]);
|
||||
expect(
|
||||
resolveGatewaySessionStoreTargetsReadOnly({
|
||||
cfg,
|
||||
targets: [{ key: "agent:retired:main" }],
|
||||
}),
|
||||
).toMatchObject([{ storePath: retiredStorePath, store: {} }]);
|
||||
expect(fs.existsSync(sqlitePath!)).toBe(false);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -3574,6 +3641,9 @@ describe("gateway session utils", () => {
|
||||
includeStoreChildEntries: true,
|
||||
}),
|
||||
).toThrow("openclaw doctor --fix");
|
||||
expect(() =>
|
||||
resolveGatewaySessionStoreTargetsReadOnly({ cfg, targets: [{ key: "main" }] }),
|
||||
).toThrow("openclaw doctor --fix");
|
||||
});
|
||||
} finally {
|
||||
resetConfigRuntimeState();
|
||||
@@ -3647,6 +3717,32 @@ describe("gateway session utils", () => {
|
||||
expect(loaded.canonicalKey).toBe("agent:main:main");
|
||||
expect(loaded.storePath).toBe(path.resolve(deletedStorePath));
|
||||
expect(loaded.entry?.sessionId).toBe("sess-deleted-main");
|
||||
closeOpenClawAgentDatabasesForTest();
|
||||
const parse = JSON.parse;
|
||||
let liveDefaultParses = 0;
|
||||
const parseSpy = vi.spyOn(JSON, "parse").mockImplementation((text, reviver) => {
|
||||
if (text.includes('"sessionId":"sess-live-default"')) {
|
||||
liveDefaultParses += 1;
|
||||
}
|
||||
return parse(text, reviver);
|
||||
});
|
||||
try {
|
||||
expect(
|
||||
resolveGatewaySessionStoreTargetsReadOnly({
|
||||
cfg,
|
||||
targets: [{ key: "agent:main:main" }],
|
||||
}),
|
||||
).toMatchObject([
|
||||
{
|
||||
agentId: "main",
|
||||
storePath: path.resolve(deletedStorePath),
|
||||
store: { "agent:main:main": { sessionId: "sess-deleted-main" } },
|
||||
},
|
||||
]);
|
||||
expect(liveDefaultParses).toBe(0);
|
||||
} finally {
|
||||
parseSpy.mockRestore();
|
||||
}
|
||||
});
|
||||
} finally {
|
||||
resetConfigRuntimeState();
|
||||
@@ -3685,6 +3781,18 @@ describe("gateway session utils", () => {
|
||||
requestedKey === key ? "incognito-owner" : undefined,
|
||||
);
|
||||
expect(target.store["agent:main:main"]).toBeUndefined();
|
||||
const [batched] = resolveGatewaySessionStoreTargetsReadOnly({
|
||||
cfg,
|
||||
targets: [{ key: requestedKey }],
|
||||
});
|
||||
expect(batched).toMatchObject({
|
||||
storePath: resolveIncognitoOpenClawAgentSqlitePath({ agentId: "main" }),
|
||||
storeKeys: [requestedKey],
|
||||
});
|
||||
expect(batched?.store[requestedKey]?.sessionId).toBe(
|
||||
requestedKey === key ? "incognito-owner" : undefined,
|
||||
);
|
||||
expect(batched?.store["agent:main:main"]).toBeUndefined();
|
||||
}
|
||||
});
|
||||
},
|
||||
@@ -3719,6 +3827,9 @@ describe("gateway session utils", () => {
|
||||
setRuntimeConfigSnapshot(cfg, cfg);
|
||||
|
||||
expect(() => loadSessionEntry("agent:main:work")).toThrow("openclaw doctor --fix");
|
||||
expect(() =>
|
||||
resolveGatewaySessionStoreTargetsReadOnly({ cfg, targets: [{ key: "agent:main:work" }] }),
|
||||
).toThrow("openclaw doctor --fix");
|
||||
});
|
||||
} finally {
|
||||
resetConfigRuntimeState();
|
||||
|
||||
@@ -1,15 +1,18 @@
|
||||
import { expectDefined } from "@openclaw/normalization-core";
|
||||
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
|
||||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import { parseAgentSessionKey } from "../routing/session-key.js";
|
||||
import type { TaskRecord } from "../tasks/task-registry.types.js";
|
||||
import { hasOperatorBoundary } from "./operator-role-policy.js";
|
||||
import type { GatewayClient } from "./server-methods/types.js";
|
||||
import type { SessionSharingTarget } from "./session-sharing-policy.js";
|
||||
import {
|
||||
authorizeIncognitoSessionTarget,
|
||||
authorizeSessionSharingTarget,
|
||||
createSessionListEntryFilter,
|
||||
isGatewayAdmin,
|
||||
resolveSessionSharingTarget,
|
||||
resolveSessionSharingTargets,
|
||||
} from "./session-sharing.js";
|
||||
|
||||
export function resolveTaskRequesterSessionTarget(
|
||||
@@ -36,7 +39,21 @@ export function canAccessTaskRequesterSession(params: {
|
||||
if (!target || isGatewayAdmin(params.client)) {
|
||||
return true;
|
||||
}
|
||||
const sharingTarget = resolveSessionSharingTarget({ cfg: params.cfg, ...target });
|
||||
return canAccessResolvedTaskSession(
|
||||
params,
|
||||
target,
|
||||
resolveSessionSharingTarget({ cfg: params.cfg, ...target }),
|
||||
);
|
||||
}
|
||||
|
||||
function canAccessResolvedTaskSession(
|
||||
params: Pick<Parameters<typeof canAccessTaskRequesterSession>[0], "cfg" | "client" | "access">,
|
||||
target: ReturnType<typeof resolveTaskRequesterSessionTarget>,
|
||||
sharingTarget: SessionSharingTarget | null,
|
||||
): boolean {
|
||||
if (!target || isGatewayAdmin(params.client)) {
|
||||
return true;
|
||||
}
|
||||
if (
|
||||
authorizeIncognitoSessionTarget({
|
||||
client: params.client,
|
||||
@@ -65,3 +82,40 @@ export function canAccessTaskRequesterSession(params: {
|
||||
});
|
||||
return visibilityFilter?.(sharingTarget.storeKey, sharingTarget.entry) ?? true;
|
||||
}
|
||||
|
||||
/** Prepare only this slice's entries; the registry drops this filter before yielding. */
|
||||
export function prepareTaskSessionReadFilter(
|
||||
params: { cfg: OpenClawConfig; client: GatewayClient | null },
|
||||
tasks: readonly Readonly<TaskRecord>[],
|
||||
): (task: Readonly<TaskRecord>) => boolean {
|
||||
if (isGatewayAdmin(params.client)) {
|
||||
return (task) => canAccessTaskRequesterSession({ ...params, task });
|
||||
}
|
||||
const requests: Array<{
|
||||
task: Readonly<TaskRecord>;
|
||||
target: ReturnType<typeof resolveTaskRequesterSessionTarget>;
|
||||
sharingTarget: SessionSharingTarget | null;
|
||||
}> = tasks.map((task) => ({
|
||||
task,
|
||||
target: resolveTaskRequesterSessionTarget(task),
|
||||
sharingTarget: null,
|
||||
}));
|
||||
const lookups = requests.flatMap((request) =>
|
||||
request.target ? [{ request, target: request.target }] : [],
|
||||
);
|
||||
for (const [index, sharingTarget] of resolveSessionSharingTargets({
|
||||
cfg: params.cfg,
|
||||
targets: lookups.map((lookup) => lookup.target),
|
||||
}).entries()) {
|
||||
expectDefined(lookups[index], "prepared task session lookup").request.sharingTarget =
|
||||
sharingTarget;
|
||||
}
|
||||
const prepared = new Map(requests.map((request) => [request.task, request]));
|
||||
return (task) => {
|
||||
const request = expectDefined(
|
||||
prepared.get(task),
|
||||
"task belongs to the synchronous access slice",
|
||||
);
|
||||
return canAccessResolvedTaskSession(params, request.target, request.sharingTarget);
|
||||
};
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import {
|
||||
listTaskRecordPage,
|
||||
resetTaskRegistryForTests,
|
||||
} from "./task-registry-query.js";
|
||||
import { markTaskTerminalById } from "./task-registry-record-api.js";
|
||||
import { configureTaskRegistryRuntime } from "./task-registry.store.js";
|
||||
import type { TaskRecord } from "./task-registry.types.js";
|
||||
|
||||
@@ -32,6 +33,59 @@ async function readTaskPage(params: Parameters<typeof listTaskRecordPage>[0]) {
|
||||
}
|
||||
|
||||
describe("listTaskRecordPage", () => {
|
||||
it.each([
|
||||
{ count: 32, mutationTurn: 1, completes: true },
|
||||
{ count: 64, mutationTurn: 2, completes: true },
|
||||
{ count: 33, mutationTurn: 1, completes: false },
|
||||
{ count: 65, mutationTurn: 2, completes: false },
|
||||
])(
|
||||
"finishes complete batches but yields unfinished work ($count tasks)",
|
||||
async ({ count, mutationTurn, completes }) => {
|
||||
configureTaskSnapshot(
|
||||
Array.from({ length: count }, (_, index): TaskRecord => ({
|
||||
taskId: `task-${index}`,
|
||||
runtime: "cli",
|
||||
requesterSessionKey: "agent:main:main",
|
||||
ownerKey: "agent:main:main",
|
||||
scopeKind: "session",
|
||||
task: "Task with queued activity",
|
||||
status: "running",
|
||||
deliveryStatus: "pending",
|
||||
notifyPolicy: "done_only",
|
||||
createdAt: 1,
|
||||
lastEventAt: 1,
|
||||
})),
|
||||
);
|
||||
let turn = 0;
|
||||
let mutations = 0;
|
||||
const update = () => {
|
||||
turn += 1;
|
||||
if (turn >= mutationTurn) {
|
||||
mutations += 1;
|
||||
markTaskTerminalById({
|
||||
taskId: "task-0",
|
||||
status: "succeeded",
|
||||
endedAt: mutations + 1,
|
||||
});
|
||||
}
|
||||
pending = setImmediate(update);
|
||||
};
|
||||
let pending = setImmediate(update);
|
||||
try {
|
||||
const page = await listTaskRecordPage({ offset: 0, limit: count });
|
||||
if (completes) {
|
||||
expect(page.ok).toBe(true);
|
||||
expect(mutations).toBe(0);
|
||||
} else {
|
||||
expect(page).toEqual({ ok: false, error: "registry_changed" });
|
||||
expect(mutations).toBeGreaterThanOrEqual(3);
|
||||
}
|
||||
} finally {
|
||||
clearImmediate(pending);
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
it("keeps large page scans responsive and sorts only the selected window", async () => {
|
||||
const total = 10_000;
|
||||
const offset = 13;
|
||||
|
||||
@@ -159,7 +159,9 @@ export async function listTaskRecordPage(params: {
|
||||
sessionKey?: string;
|
||||
sessionAgentId?: string;
|
||||
cfg?: OpenClawConfig;
|
||||
filter?: (task: Readonly<TaskRecord>) => boolean;
|
||||
prepareFilter?: (
|
||||
tasks: readonly Readonly<TaskRecord>[],
|
||||
) => (task: Readonly<TaskRecord>) => boolean;
|
||||
sortBy?: "updatedAt" | "endedAt";
|
||||
}): Promise<
|
||||
Result<
|
||||
@@ -186,40 +188,48 @@ export async function listTaskRecordPage(params: {
|
||||
let matchingCount = 0;
|
||||
let heapReady = false;
|
||||
let scannedCount = 0;
|
||||
for (const task of tasks.values()) {
|
||||
if (scannedCount >= scanLimit) {
|
||||
break;
|
||||
}
|
||||
scannedCount += 1;
|
||||
// Yield large scans in small deterministic slices so task history cannot
|
||||
// monopolize the Gateway event loop while other requests are waiting.
|
||||
if (scannedCount % 32 === 0) {
|
||||
const iterator = tasks.values();
|
||||
let current = iterator.next();
|
||||
while (!current.done && scannedCount < scanLimit) {
|
||||
// Yield only when another batch exists; completed pages keep their revision.
|
||||
if (scannedCount > 0) {
|
||||
await yieldToEventLoop();
|
||||
}
|
||||
if (
|
||||
(statuses && !statuses.has(task.status)) ||
|
||||
!taskMatchesAgent(task, agentId, params.cfg) ||
|
||||
!taskMatchesRelatedSession(task, sessionKey, params.sessionAgentId, params.cfg) ||
|
||||
(params.filter && !params.filter(task))
|
||||
) {
|
||||
continue;
|
||||
const batch: TaskRecord[] = [];
|
||||
while (!current.done && batch.length < 32 && scannedCount < scanLimit) {
|
||||
batch.push(current.value);
|
||||
scannedCount += 1;
|
||||
current = iterator.next();
|
||||
}
|
||||
matchingCount += 1;
|
||||
if (windowSize <= 0) {
|
||||
continue;
|
||||
}
|
||||
if (window.length < windowSize) {
|
||||
window.push(task);
|
||||
continue;
|
||||
}
|
||||
if (!heapReady) {
|
||||
heapifyWorstTaskFirst(window, compare);
|
||||
heapReady = true;
|
||||
}
|
||||
const cutoff = window[0];
|
||||
if (cutoff && compare(task, cutoff) < 0) {
|
||||
window[0] = task;
|
||||
siftWorstTaskDown(window, 0, compare);
|
||||
const candidates = batch.filter(
|
||||
(task) =>
|
||||
(!statuses || statuses.has(task.status)) &&
|
||||
taskMatchesAgent(task, agentId, params.cfg) &&
|
||||
taskMatchesRelatedSession(task, sessionKey, params.sessionAgentId, params.cfg),
|
||||
);
|
||||
// Prepared metadata belongs to this synchronous slice, never the next await.
|
||||
const filter = params.prepareFilter?.(candidates);
|
||||
for (const task of candidates) {
|
||||
if (filter && !filter(task)) {
|
||||
continue;
|
||||
}
|
||||
matchingCount += 1;
|
||||
if (windowSize <= 0) {
|
||||
continue;
|
||||
}
|
||||
if (window.length < windowSize) {
|
||||
window.push(task);
|
||||
continue;
|
||||
}
|
||||
if (!heapReady) {
|
||||
heapifyWorstTaskFirst(window, compare);
|
||||
heapReady = true;
|
||||
}
|
||||
const cutoff = window[0];
|
||||
if (cutoff && compare(task, cutoff) < 0) {
|
||||
window[0] = task;
|
||||
siftWorstTaskDown(window, 0, compare);
|
||||
}
|
||||
}
|
||||
}
|
||||
if (revision !== readTaskRegistryRevision()) {
|
||||
|
||||
Reference in New Issue
Block a user