Skip to content

Commit 0bf66ab

Browse files
committed
fix(sessions): scope ambient transcript watermark to session id
Ambient transcript watermarks now carry the transcript session id, resolve only for the current session entry, and skip stale room-event hooks that no longer match the prepared transcript session. This protects Telegram group prompt windows after reset by backfilling rows that are no longer present in the new session transcript, while preserving steady-state watermark filtering within one session. Fixes #99373 Release-note: fixes Telegram group context loss after session reset when ambient transcript watermarks outlived the transcript they referenced.
1 parent 08079ec commit 0bf66ab

7 files changed

Lines changed: 265 additions & 3 deletions

File tree

extensions/telegram/src/bot-message-context.prompt-context.test.ts

Lines changed: 118 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,14 @@
1-
import { describe, expect, it } from "vitest";
1+
import fs from "node:fs";
2+
import os from "node:os";
3+
import path from "node:path";
4+
import {
5+
getSessionEntry,
6+
readAmbientTranscriptWatermark,
7+
resolveAmbientTranscriptWatermarkKey,
8+
updateAmbientTranscriptWatermark,
9+
upsertSessionEntry,
10+
} from "openclaw/plugin-sdk/session-store-runtime";
11+
import { afterEach, describe, expect, it } from "vitest";
212
import { buildTelegramMessageContextForTest } from "./bot-message-context.test-harness.js";
313
import type { TelegramPromptContextEntry } from "./bot-message-context.types.js";
414

@@ -20,6 +30,20 @@ const telegramChatWindowContext: TelegramPromptContextEntry = {
2030
},
2131
};
2232

33+
const tempDirs: string[] = [];
34+
35+
function createTempSessionStorePath(): string {
36+
const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-telegram-watermark-"));
37+
tempDirs.push(tempDir);
38+
return path.join(tempDir, "sessions.json");
39+
}
40+
41+
afterEach(() => {
42+
for (const tempDir of tempDirs.splice(0)) {
43+
fs.rmSync(tempDir, { recursive: true, force: true });
44+
}
45+
});
46+
2347
describe("buildTelegramMessageContext prompt context", () => {
2448
it("omits Telegram chat-window context for existing unthreaded private DM sessions", async () => {
2549
const ctx = await buildTelegramMessageContextForTest({
@@ -180,6 +204,7 @@ describe("buildTelegramMessageContext prompt context", () => {
180204
readAmbientTranscriptWatermark: ({ key }) =>
181205
key === '["telegram","default","-1001234567890",""]'
182206
? {
207+
sessionId: "session-current",
183208
messageId: "11",
184209
timestampMs: 1_700_000_001_000,
185210
updatedAt: 1_700_000_003_000,
@@ -237,6 +262,7 @@ describe("buildTelegramMessageContext prompt context", () => {
237262
]),
238263
sessionRuntime: {
239264
readAmbientTranscriptWatermark: () => ({
265+
sessionId: "session-current",
240266
messageId: "11",
241267
timestampMs: 1_700_000_001_000,
242268
updatedAt: 1_700_000_003_000,
@@ -286,6 +312,7 @@ describe("buildTelegramMessageContext prompt context", () => {
286312
readAmbientTranscriptWatermark: ({ key }) =>
287313
key === '["telegram","default","-1001234567890",""]'
288314
? {
315+
sessionId: "session-current",
289316
messageId: "11",
290317
timestampMs: 1_700_000_001_000,
291318
updatedAt: 1_700_000_003_000,
@@ -306,4 +333,94 @@ describe("buildTelegramMessageContext prompt context", () => {
306333
expect(ctx.ctxPayload.InboundHistory).toBeUndefined();
307334
expect(ctx.ctxPayload.UntrustedStructuredContext).toBeUndefined();
308335
});
336+
337+
it("backfills Telegram group history when the ambient watermark belongs to a reset session", async () => {
338+
const storePath = createTempSessionStorePath();
339+
const sessionKey = "agent:main:telegram:group:-1001234567890";
340+
const key = resolveAmbientTranscriptWatermarkKey({
341+
channel: "telegram",
342+
accountId: "default",
343+
conversationId: "-1001234567890",
344+
});
345+
346+
await upsertSessionEntry({
347+
storePath,
348+
sessionKey,
349+
entry: { sessionId: "before-reset", updatedAt: 1_700_000_000_000 },
350+
});
351+
await updateAmbientTranscriptWatermark({
352+
storePath,
353+
sessionKey,
354+
key,
355+
messageId: "11",
356+
timestampMs: 1_700_000_001_000,
357+
});
358+
const persistedEntry = getSessionEntry({ storePath, sessionKey });
359+
if (!persistedEntry) {
360+
throw new Error("Expected persisted session entry");
361+
}
362+
await upsertSessionEntry({
363+
storePath,
364+
sessionKey,
365+
entry: {
366+
...persistedEntry,
367+
sessionId: "after-reset",
368+
updatedAt: 1_700_000_002_000,
369+
},
370+
});
371+
372+
const ctx = await buildTelegramMessageContextForTest({
373+
message: {
374+
message_id: 13,
375+
chat: { id: -1001234567890, type: "supergroup", title: "Forum" },
376+
from: { id: 1234, first_name: "Pat" },
377+
text: "@bot what happened?",
378+
entities: [{ type: "mention", offset: 0, length: 4 }],
379+
},
380+
historyLimit: 10,
381+
groupHistories: new Map([
382+
[
383+
"-1001234567890",
384+
[
385+
{
386+
messageId: "10",
387+
sender: "Sam",
388+
timestamp: 1_700_000_000_000,
389+
body: "persisted ambient one",
390+
},
391+
{
392+
messageId: "11",
393+
sender: "Lee",
394+
timestamp: 1_700_000_001_000,
395+
body: "persisted ambient two",
396+
},
397+
{
398+
messageId: "12",
399+
sender: "Mira",
400+
timestamp: 1_700_000_002_000,
401+
body: "unpersisted gap",
402+
},
403+
],
404+
],
405+
]),
406+
sessionRuntime: {
407+
readAmbientTranscriptWatermark,
408+
resolveAmbientTranscriptWatermarkKey,
409+
resolveStorePath: () => storePath,
410+
},
411+
});
412+
413+
expect(ctx?.ctxPayload.UntrustedStructuredContext).toEqual([
414+
expect.objectContaining({
415+
type: "chat_window",
416+
payload: expect.objectContaining({
417+
messages: [
418+
expect.objectContaining({ message_id: "10", body: "persisted ambient one" }),
419+
expect.objectContaining({ message_id: "11", body: "persisted ambient two" }),
420+
expect.objectContaining({ message_id: "12", body: "unpersisted gap" }),
421+
],
422+
}),
423+
}),
424+
]);
425+
});
309426
});

extensions/telegram/src/bot.test.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2107,6 +2107,7 @@ describe("createTelegramBot", () => {
21072107
() => "telegram:default:42",
21082108
);
21092109
telegramBotDepsForTest.readAmbientTranscriptWatermark = vi.fn(() => ({
2110+
sessionId: "session-current",
21102111
messageId: "502",
21112112
timestampMs: 1_736_380_860_000,
21122113
updatedAt: 1_736_380_900_000,

src/auto-reply/reply/get-reply-run.media-only.test.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2055,6 +2055,7 @@ describe("runPreparedReply media-only handling", () => {
20552055
key: '["telegram","","-100123",""]',
20562056
messageId: "35676",
20572057
timestampMs: 1_710_000_000_000,
2058+
expectedSessionId: expect.any(String),
20582059
});
20592060
expect(call?.followupRun.currentInboundContext?.text).toContain(
20602061
"#35675 obviyus ->#35674: Are you fr fr",

src/auto-reply/reply/get-reply-run.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -174,6 +174,7 @@ function normalizeMessageTimestampMs(value: unknown): number | undefined {
174174
}
175175

176176
async function updateRoomEventAmbientTranscriptWatermark(params: {
177+
expectedSessionId: string;
177178
sessionCtx: TemplateContext;
178179
storePath?: string;
179180
sessionKey?: string;
@@ -191,6 +192,7 @@ async function updateRoomEventAmbientTranscriptWatermark(params: {
191192
key,
192193
messageId,
193194
timestampMs: params.sessionCtx.AmbientTranscriptTimestampMs,
195+
expectedSessionId: params.expectedSessionId,
194196
});
195197
}
196198

@@ -1307,6 +1309,7 @@ export async function runPreparedReply(
13071309
onMessagePersisted: isRoomEvent
13081310
? async () =>
13091311
await updateRoomEventAmbientTranscriptWatermark({
1312+
expectedSessionId: preparedSessionState.sessionId,
13101313
sessionCtx,
13111314
storePath,
13121315
sessionKey: sessionKey ?? preparedSessionState.sessionId,
Lines changed: 124 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,124 @@
1+
import fs from "node:fs";
2+
import os from "node:os";
3+
import path from "node:path";
4+
import { afterEach, beforeEach, describe, expect, it } from "vitest";
5+
import {
6+
readAmbientTranscriptWatermark,
7+
resolveAmbientTranscriptWatermarkKey,
8+
updateAmbientTranscriptWatermark,
9+
} from "./ambient-transcript-watermark.js";
10+
import { loadSessionEntry, replaceSessionEntry } from "./session-accessor.js";
11+
12+
describe("ambient transcript watermark", () => {
13+
let tempDir: string;
14+
let storePath: string;
15+
const sessionKey = "agent:main:telegram:group:-100123";
16+
const key = resolveAmbientTranscriptWatermarkKey({
17+
channel: "telegram",
18+
accountId: "default",
19+
conversationId: "-100123",
20+
});
21+
22+
beforeEach(() => {
23+
tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-ambient-watermark-"));
24+
storePath = path.join(tempDir, "sessions.json");
25+
});
26+
27+
afterEach(() => {
28+
fs.rmSync(tempDir, { recursive: true, force: true });
29+
});
30+
31+
it("stamps and resolves the watermark for the current session id only", async () => {
32+
await replaceSessionEntry(
33+
{ sessionKey, storePath },
34+
{ sessionId: "before-reset", updatedAt: 1_700_000_000_000 },
35+
);
36+
37+
await updateAmbientTranscriptWatermark({
38+
storePath,
39+
sessionKey,
40+
key,
41+
messageId: "11",
42+
timestampMs: 1_700_000_001_000,
43+
});
44+
45+
const persistedEntry = loadSessionEntry({ sessionKey, storePath });
46+
if (!persistedEntry) {
47+
throw new Error("Expected persisted session entry");
48+
}
49+
expect(persistedEntry?.ambientTranscriptWatermarks?.[key]).toMatchObject({
50+
sessionId: "before-reset",
51+
messageId: "11",
52+
timestampMs: 1_700_000_001_000,
53+
});
54+
expect(readAmbientTranscriptWatermark(persistedEntry, key)).toMatchObject({
55+
sessionId: "before-reset",
56+
messageId: "11",
57+
});
58+
59+
await replaceSessionEntry(
60+
{ sessionKey, storePath },
61+
{
62+
...persistedEntry,
63+
sessionId: "after-reset",
64+
updatedAt: 1_700_000_002_000,
65+
},
66+
);
67+
68+
const resetEntry = loadSessionEntry({ sessionKey, storePath });
69+
expect(readAmbientTranscriptWatermark(resetEntry, key)).toBeUndefined();
70+
71+
await updateAmbientTranscriptWatermark({
72+
storePath,
73+
sessionKey,
74+
key,
75+
messageId: "12",
76+
timestampMs: 1_700_000_002_000,
77+
expectedSessionId: "before-reset",
78+
});
79+
80+
expect(
81+
readAmbientTranscriptWatermark(loadSessionEntry({ sessionKey, storePath }), key),
82+
).toBeUndefined();
83+
84+
await updateAmbientTranscriptWatermark({
85+
storePath,
86+
sessionKey,
87+
key,
88+
messageId: "12",
89+
timestampMs: 1_700_000_002_000,
90+
expectedSessionId: "after-reset",
91+
});
92+
93+
expect(
94+
readAmbientTranscriptWatermark(loadSessionEntry({ sessionKey, storePath }), key),
95+
).toMatchObject({
96+
sessionId: "after-reset",
97+
messageId: "12",
98+
});
99+
});
100+
101+
it("ignores legacy watermarks without a session id", () => {
102+
fs.writeFileSync(
103+
storePath,
104+
JSON.stringify({
105+
[sessionKey]: {
106+
sessionId: "current-session",
107+
updatedAt: 1_700_000_000_000,
108+
ambientTranscriptWatermarks: {
109+
[key]: {
110+
messageId: "11",
111+
timestampMs: 1_700_000_001_000,
112+
updatedAt: 1_700_000_002_000,
113+
},
114+
},
115+
},
116+
}),
117+
"utf-8",
118+
);
119+
120+
expect(
121+
readAmbientTranscriptWatermark(loadSessionEntry({ sessionKey, storePath }), key),
122+
).toBeUndefined();
123+
});
124+
});

src/config/sessions/ambient-transcript-watermark.ts

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -52,10 +52,14 @@ function isAmbientTranscriptWatermarkAfter(
5252
}
5353

5454
export function readAmbientTranscriptWatermark(
55-
entry: Pick<SessionEntry, "ambientTranscriptWatermarks"> | undefined,
55+
entry: Pick<SessionEntry, "ambientTranscriptWatermarks" | "sessionId"> | undefined,
5656
key: string,
5757
): AmbientTranscriptWatermark | undefined {
58-
return entry?.ambientTranscriptWatermarks?.[key];
58+
const watermark = entry?.ambientTranscriptWatermarks?.[key];
59+
// A watermark only vouches for rows in the transcript it was written against.
60+
// After a session reset those rows live in an archived file the model never
61+
// reads, so a cross-session (or legacy sessionId-less) watermark must not hide them.
62+
return watermark?.sessionId === entry?.sessionId ? watermark : undefined;
5963
}
6064

6165
export async function updateAmbientTranscriptWatermark(params: {
@@ -64,13 +68,23 @@ export async function updateAmbientTranscriptWatermark(params: {
6468
key: string;
6569
messageId: string;
6670
timestampMs?: number;
71+
expectedSessionId?: string;
6772
}): Promise<SessionEntry | null> {
6873
return await updateSessionEntry(
6974
{
7075
storePath: params.storePath,
7176
sessionKey: params.sessionKey,
7277
},
7378
(entry) => {
79+
// onMessagePersisted fires after the durable row write; if the session was
80+
// reset in between, stamping the new sessionId would hide rows that only
81+
// exist in the archived transcript. Skip the advance instead.
82+
if (!entry.sessionId) {
83+
return null;
84+
}
85+
if (params.expectedSessionId !== undefined && entry.sessionId !== params.expectedSessionId) {
86+
return null;
87+
}
7488
const current = readAmbientTranscriptWatermark(entry, params.key);
7589
if (
7690
!isAmbientTranscriptWatermarkAfter(
@@ -84,6 +98,7 @@ export async function updateAmbientTranscriptWatermark(params: {
8498
ambientTranscriptWatermarks: {
8599
...entry.ambientTranscriptWatermarks,
86100
[params.key]: {
101+
sessionId: entry.sessionId,
87102
messageId: params.messageId,
88103
...(params.timestampMs !== undefined ? { timestampMs: params.timestampMs } : {}),
89104
updatedAt: Date.now(),

src/config/sessions/types.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -113,6 +113,7 @@ export type SessionContextBudgetStatus = {
113113
};
114114

115115
export type AmbientTranscriptWatermark = {
116+
sessionId: string;
116117
messageId: string;
117118
timestampMs?: number;
118119
updatedAt: number;

0 commit comments

Comments
 (0)