fix(sessions): route snapshot transcript writes by owner

This commit is contained in:
Ayaan Zaidi
2026-08-12 19:40:31 +02:00
parent 7d1bd1ad6e
commit 5477cc757f
18 changed files with 117 additions and 39 deletions
@@ -1 +1 @@
{"contentHash":"48e9f1762613441ec71879de8fcfffcf0d382867b55436716ee00459d7b4f49b","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"}
{"contentHash":"8f18517eea8e82b6fc749ae5cb42d21b4f1406b51e16c70203205315983e216f","entrypoint":"agent-harness-runtime","importSpecifier":"openclaw/plugin-sdk/agent-harness-runtime"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"fbc1f78eb076fb8e8e966f7cbf377303ec162cc3f2c5fa50e146d2e73618560f","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"}
{"contentHash":"13aa9f6d8fee157fdf4d98d524da8e21423eaad57ec3612ee98df483ad8bd507","entrypoint":"agent-harness","importSpecifier":"openclaw/plugin-sdk/agent-harness"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"02987e8974ddb4d4ac590a9d9fedb620689a7ef198b78c4cb87120dbb546fbd2","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"}
{"contentHash":"6d225dda2653c8558cc42665fb3ba787fd4dd02f970b76e8eb2386b053dab6cf","entrypoint":"channel-core","importSpecifier":"openclaw/plugin-sdk/channel-core"}
@@ -1 +1 @@
{"contentHash":"a531034eea4aa11ef247468055c90ac2be9b90eed15d20bcd621ba89647e61b3","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"}
{"contentHash":"2a322de029d348559b6e25414b9ca484fb440ab6a832440bb7b9744eefa25a03","entrypoint":"channel-entry-contract","importSpecifier":"openclaw/plugin-sdk/channel-entry-contract"}
@@ -1 +1 @@
{"contentHash":"db916c141015c7ec52d0f3efc7fcb855084e81d8502d610d63f3eb78f5b7a8e5","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"}
{"contentHash":"d4e97e41bf610c01a2bd48308a2a60f8951cd9c5d204b8575f43515cf6e897bd","entrypoint":"channel-message","importSpecifier":"openclaw/plugin-sdk/channel-message"}
@@ -1 +1 @@
{"contentHash":"7ac53597410b8bd7a525462e3e1598408010711d9cb2b77d68ad69332059ba7a","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"}
{"contentHash":"ec6e168cb49b72994adc5ee6e5eaae7c2181cf51f2bb61ce74d6f8d186d0148d","entrypoint":"channel-outbound","importSpecifier":"openclaw/plugin-sdk/channel-outbound"}
@@ -1 +1 @@
{"contentHash":"7cf6e6ddd0bbb3c40fb012760ebcdce17834d9347bef430d4e7f2efd93db1f1e","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"}
{"contentHash":"60290b8598657d224562ca9eb78377193da6e91da05b0373240af43e05e1b241","entrypoint":"channel-plugin-common","importSpecifier":"openclaw/plugin-sdk/channel-plugin-common"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"0c76b2fd651acded8caf3f8c76f7ee292b530b6e0a052cb4f0688fc0bc00a1c0","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"}
{"contentHash":"ed28932f6d3093ad7399064865a14afd44e7a6272d957ec0f919786f0dc12f62","entrypoint":"core","importSpecifier":"openclaw/plugin-sdk/core"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"19c444b6b2eb6062adabf88dc92060553042071de32d1113e2dae4a528d24fa6","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"}
{"contentHash":"9bca101dfadd456fa88d33bb610d2f8290b840110bddfbeac3bbe32eb8677edd","entrypoint":"discord","importSpecifier":"openclaw/plugin-sdk/discord"}
@@ -1 +1 @@
{"contentHash":"ae337f8842b30f3d63ae5dd90963e7b73894cff02a3aea4b6f218ed0715fdd74","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"}
{"contentHash":"0c1ea4c591fa0c0281bbe7bf261f9f4e39977b52133763c95a5cffa91c47b89d","entrypoint":"inbound-reply-dispatch","importSpecifier":"openclaw/plugin-sdk/inbound-reply-dispatch"}
@@ -1 +1 @@
{"contentHash":"8c90eea3929b848916e5d4aaa67e8dac636e47e8ab766a5260b888d3f2dd6cfb","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"}
{"contentHash":"91b285baa5bfa0a040cb54e63971e4bcebcb14b11ca206ec7d59cf20cd7f04be","entrypoint":"meeting-runtime","importSpecifier":"openclaw/plugin-sdk/meeting-runtime"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"ec039f279134221e3937ed3388e0cc0f065e2ab0a62a80b42384d7a10cbd98c3","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"}
{"contentHash":"ffc87bb6abfe73b76cb379a9d6da485cce98500053f4ca2b3abf338ef8178e95","entrypoint":"plugin-entry","importSpecifier":"openclaw/plugin-sdk/plugin-entry"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"c86462cf73045c6a2eb37a77be545b3148b58890dd0b112750d24b0829ad6166","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"}
{"contentHash":"470acf51a08826d82d08316d0f9607e97d954ffd840aafa9bb6bd9aef0106ba7","entrypoint":"plugin-runtime","importSpecifier":"openclaw/plugin-sdk/plugin-runtime"}
@@ -1 +1 @@
{"contentHash":"9cd0143763f80869c7182fcab934d493a0106d0b358d989fd64270fae540c73c","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"}
{"contentHash":"5003150fed73cdbb9a19e4309a070fa657c7625e1b44b36d44f987131379afde","entrypoint":"provider-catalog-runtime","importSpecifier":"openclaw/plugin-sdk/provider-catalog-runtime"}
+1 -1
View File
@@ -1 +1 @@
{"contentHash":"90cb8357bc1472f7b5df268aa40850fbc8d33e401e3f9d5b7ea73f7b0842b0ea","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"}
{"contentHash":"026a053659a9f728a9452da4d69f4581213074eeea930f07000377cbac48a9d7","entrypoint":"tool-plugin","importSpecifier":"openclaw/plugin-sdk/tool-plugin"}
@@ -1 +1 @@
{"contentHash":"81022c7a14e82a7d0b50aaec10fa9d6c5494e2cd8398aa1bc4fe4556f5fbfdcf","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"}
{"contentHash":"db7793d03f781b3044dce9497370f591b44fdcec2afc36d02e173e712c6e66fd","entrypoint":"webhook-ingress","importSpecifier":"openclaw/plugin-sdk/webhook-ingress"}
@@ -281,8 +281,8 @@ async function persistExpectedSessionTranscriptTurn(
expectedSessionId: string;
},
): Promise<SessionTranscriptTurnPersistResult> {
const sessionKey = scope.sessionKey?.trim();
if (!scope.storePath || !sessionKey) {
const requestedSessionKey = scope.sessionKey?.trim();
if (!scope.storePath || !requestedSessionKey) {
throw new Error("Cannot guard a transcript turn without a session store and key");
}
const storePath = scope.storePath;
@@ -290,11 +290,19 @@ async function persistExpectedSessionTranscriptTurn(
const agentId = resolveTranscriptTurnAgentId({
config: options.config ?? getRuntimeConfig(),
scopeAgentId: scope.agentId,
sessionKey,
sessionKey: requestedSessionKey,
storePath,
sessionStore: scope.sessionStore,
env: scope.env,
});
const runtimeTarget = await resolveSessionTranscriptRuntimeTarget({
agentId,
...(scope.env ? { env: scope.env } : {}),
sessionId: expectedSessionId,
sessionKey: requestedSessionKey,
storePath,
});
const sessionKey = runtimeTarget.sessionKey;
const resolved = scope.sessionStore
? resolveSessionEntryFromStore({ store: scope.sessionStore, sessionKey })
: resolveSessionEntrySelection({
@@ -355,8 +363,13 @@ async function persistExpectedSessionTranscriptTurn(
};
}
// The requested key remains the caller's live update route; the resolved
// target above is the distinct physical SQLite owner.
await publishTranscriptTurnUpdate({
target,
target:
requestedSessionKey === target.sessionKey
? target
: { ...target, sessionKey: requestedSessionKey },
sessionEntry: turn.sessionEntry,
updateMode: options.updateMode ?? "inline",
publishWhen: options.publishWhen ?? "when-appended",
@@ -404,16 +417,17 @@ async function resolveTranscriptTurnTarget(
agentId,
env: scope.env,
});
const runtimeTarget = scope.sessionStore
? undefined
: await resolveSessionTranscriptRuntimeTarget({
agentId,
...(scope.env ? { env: scope.env } : {}),
sessionId: scope.sessionId,
sessionKey,
storePath,
});
const resolvedSessionKey = runtimeTarget?.sessionKey ?? sessionKey;
// A caller snapshot may retain the routing key that admitted the turn. The
// persisted window owns durable writes; resolving it is read-only, so a
// memory-only mirror still avoids materializing SQLite state.
const runtimeTarget = await resolveSessionTranscriptRuntimeTarget({
agentId,
...(scope.env ? { env: scope.env } : {}),
sessionId: scope.sessionId,
sessionKey,
storePath,
});
const resolvedSessionKey = runtimeTarget.sessionKey;
const resolved = scope.sessionStore
? resolveSessionEntryFromStore({ store: scope.sessionStore, sessionKey: resolvedSessionKey })
: undefined;
@@ -425,11 +439,11 @@ async function resolveTranscriptTurnTarget(
sessionKey: resolvedSessionKey,
storePath,
});
const sessionEntry = resolved?.existing ?? scope.sessionEntry ?? persistedEntry;
const sessionEntry = persistedEntry ?? resolved?.existing ?? scope.sessionEntry;
return {
agentId,
sessionId: scope.sessionId,
sessionKey: runtimeTarget?.sessionKey ?? resolved?.normalizedKey ?? resolvedSessionKey,
sessionKey: runtimeTarget.sessionKey,
storePath,
sessionEntry,
entryFromPersistedStore: persistedEntry != null,
@@ -673,7 +673,13 @@ function createFixturePaths(prefix: string): { dir: string; transcriptPath: stri
return { dir, transcriptPath };
}
async function createTranscriptFixture(prefix: string) {
async function createTranscriptFixture(
prefix: string,
owner: Pick<SessionAccessScope, "agentId" | "sessionKey"> = {
agentId: "main",
sessionKey: "main",
},
) {
const { dir, transcriptPath } = createFixturePaths(prefix);
fs.writeFileSync(
transcriptPath,
@@ -688,11 +694,14 @@ async function createTranscriptFixture(prefix: string) {
);
// The accessor resolves transcript targets from the persisted store, so the
// fixture seeds a real entry instead of relying on the mocked gateway wrapper.
await replaceSessionEntry(sessionEntryScope(), {
sessionId: mockState.sessionId,
sessionFile: transcriptPath,
updatedAt: Date.now(),
});
await replaceSessionEntry(
{ ...owner, storePath: mockState.storePath },
{
sessionId: mockState.sessionId,
sessionFile: transcriptPath,
updatedAt: Date.now(),
},
);
return dir;
}
@@ -5040,7 +5049,10 @@ describe("chat directive tag stripping for non-streaming final payloads", () =>
});
it("chat.inject scopes selected-agent global sessions before appending", async () => {
await createTranscriptFixture("openclaw-chat-inject-selected-global-");
await createTranscriptFixture("openclaw-chat-inject-selected-global-", {
agentId: "work",
sessionKey: "agent:work:global",
});
mockState.config = {
agents: { list: [{ id: "main", default: true }, { id: "work" }] },
session: { scope: "global" },
@@ -7079,6 +7091,58 @@ describe("chat directive tag stripping for non-streaming final payloads", () =>
expect(getTotalPendingReplies()).toBe(0);
});
it("persists a Gateway user turn under the durable owner when its loaded key is stale", async () => {
createFixturePaths("openclaw-chat-send-stale-transcript-owner-");
const canonicalSessionKey = "agent:main:canonical-transcript-owner";
const staleSessionKey = "agent:main:stale-transcript-owner";
await replaceSessionEntry(
{
agentId: "main",
sessionKey: canonicalSessionKey,
storePath: mockState.storePath,
},
{ sessionId: mockState.sessionId, updatedAt: 1 },
);
mockState.finalText = "ok";
const { send } = createChatRequestFixture();
await send({
idempotencyKey: "idem-stale-transcript-owner",
message: "keep this Gateway turn",
sessionKey: staleSessionKey,
expectBroadcast: false,
});
const persistedEvents = loadTranscriptEventsSync({
agentId: "main",
sessionId: mockState.sessionId,
sessionKey: canonicalSessionKey,
storePath: mockState.storePath,
});
expect(persistedEvents).toContainEqual(
expect.objectContaining({
type: "message",
message: expect.objectContaining({
role: "user",
content: "keep this Gateway turn",
}),
}),
);
expect(
loadSqliteSessionEntry({
agentId: "main",
sessionKey: staleSessionKey,
storePath: mockState.storePath,
}),
).toBeUndefined();
expect(findUserUpdate()?.target).toEqual({
agentId: "main",
sessionId: mockState.sessionId,
sessionKey: staleSessionKey,
storePath: mockState.storePath,
});
});
it("emits a user transcript update when chat.send fails before an agent run starts", async () => {
await createGatewayUserTurnSqliteFixture("openclaw-chat-send-user-transcript-error-no-run-");
mockState.dispatchError = new Error("upstream unavailable");