From 0bf66ab7bdbf0a2a35268feee7b20036c688f3eb Mon Sep 17 00:00:00 2001 From: Ayaan Zaidi Date: Thu, 2 Jul 2026 22:40:33 -0700 Subject: [PATCH] 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. --- ...bot-message-context.prompt-context.test.ts | 119 ++++++++++++++++- extensions/telegram/src/bot.test.ts | 1 + .../reply/get-reply-run.media-only.test.ts | 1 + src/auto-reply/reply/get-reply-run.ts | 3 + .../ambient-transcript-watermark.test.ts | 124 ++++++++++++++++++ .../sessions/ambient-transcript-watermark.ts | 19 ++- src/config/sessions/types.ts | 1 + 7 files changed, 265 insertions(+), 3 deletions(-) create mode 100644 src/config/sessions/ambient-transcript-watermark.test.ts diff --git a/extensions/telegram/src/bot-message-context.prompt-context.test.ts b/extensions/telegram/src/bot-message-context.prompt-context.test.ts index fe5efab78886..519b03381a42 100644 --- a/extensions/telegram/src/bot-message-context.prompt-context.test.ts +++ b/extensions/telegram/src/bot-message-context.prompt-context.test.ts @@ -1,4 +1,14 @@ -import { describe, expect, it } from "vitest"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { + getSessionEntry, + readAmbientTranscriptWatermark, + resolveAmbientTranscriptWatermarkKey, + updateAmbientTranscriptWatermark, + upsertSessionEntry, +} from "openclaw/plugin-sdk/session-store-runtime"; +import { afterEach, describe, expect, it } from "vitest"; import { buildTelegramMessageContextForTest } from "./bot-message-context.test-harness.js"; import type { TelegramPromptContextEntry } from "./bot-message-context.types.js"; @@ -20,6 +30,20 @@ const telegramChatWindowContext: TelegramPromptContextEntry = { }, }; +const tempDirs: string[] = []; + +function createTempSessionStorePath(): string { + const tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-telegram-watermark-")); + tempDirs.push(tempDir); + return path.join(tempDir, "sessions.json"); +} + +afterEach(() => { + for (const tempDir of tempDirs.splice(0)) { + fs.rmSync(tempDir, { recursive: true, force: true }); + } +}); + describe("buildTelegramMessageContext prompt context", () => { it("omits Telegram chat-window context for existing unthreaded private DM sessions", async () => { const ctx = await buildTelegramMessageContextForTest({ @@ -180,6 +204,7 @@ describe("buildTelegramMessageContext prompt context", () => { readAmbientTranscriptWatermark: ({ key }) => key === '["telegram","default","-1001234567890",""]' ? { + sessionId: "session-current", messageId: "11", timestampMs: 1_700_000_001_000, updatedAt: 1_700_000_003_000, @@ -237,6 +262,7 @@ describe("buildTelegramMessageContext prompt context", () => { ]), sessionRuntime: { readAmbientTranscriptWatermark: () => ({ + sessionId: "session-current", messageId: "11", timestampMs: 1_700_000_001_000, updatedAt: 1_700_000_003_000, @@ -286,6 +312,7 @@ describe("buildTelegramMessageContext prompt context", () => { readAmbientTranscriptWatermark: ({ key }) => key === '["telegram","default","-1001234567890",""]' ? { + sessionId: "session-current", messageId: "11", timestampMs: 1_700_000_001_000, updatedAt: 1_700_000_003_000, @@ -306,4 +333,94 @@ describe("buildTelegramMessageContext prompt context", () => { expect(ctx.ctxPayload.InboundHistory).toBeUndefined(); expect(ctx.ctxPayload.UntrustedStructuredContext).toBeUndefined(); }); + + it("backfills Telegram group history when the ambient watermark belongs to a reset session", async () => { + const storePath = createTempSessionStorePath(); + const sessionKey = "agent:main:telegram:group:-1001234567890"; + const key = resolveAmbientTranscriptWatermarkKey({ + channel: "telegram", + accountId: "default", + conversationId: "-1001234567890", + }); + + await upsertSessionEntry({ + storePath, + sessionKey, + entry: { sessionId: "before-reset", updatedAt: 1_700_000_000_000 }, + }); + await updateAmbientTranscriptWatermark({ + storePath, + sessionKey, + key, + messageId: "11", + timestampMs: 1_700_000_001_000, + }); + const persistedEntry = getSessionEntry({ storePath, sessionKey }); + if (!persistedEntry) { + throw new Error("Expected persisted session entry"); + } + await upsertSessionEntry({ + storePath, + sessionKey, + entry: { + ...persistedEntry, + sessionId: "after-reset", + updatedAt: 1_700_000_002_000, + }, + }); + + const ctx = await buildTelegramMessageContextForTest({ + message: { + message_id: 13, + chat: { id: -1001234567890, type: "supergroup", title: "Forum" }, + from: { id: 1234, first_name: "Pat" }, + text: "@bot what happened?", + entities: [{ type: "mention", offset: 0, length: 4 }], + }, + historyLimit: 10, + groupHistories: new Map([ + [ + "-1001234567890", + [ + { + messageId: "10", + sender: "Sam", + timestamp: 1_700_000_000_000, + body: "persisted ambient one", + }, + { + messageId: "11", + sender: "Lee", + timestamp: 1_700_000_001_000, + body: "persisted ambient two", + }, + { + messageId: "12", + sender: "Mira", + timestamp: 1_700_000_002_000, + body: "unpersisted gap", + }, + ], + ], + ]), + sessionRuntime: { + readAmbientTranscriptWatermark, + resolveAmbientTranscriptWatermarkKey, + resolveStorePath: () => storePath, + }, + }); + + expect(ctx?.ctxPayload.UntrustedStructuredContext).toEqual([ + expect.objectContaining({ + type: "chat_window", + payload: expect.objectContaining({ + messages: [ + expect.objectContaining({ message_id: "10", body: "persisted ambient one" }), + expect.objectContaining({ message_id: "11", body: "persisted ambient two" }), + expect.objectContaining({ message_id: "12", body: "unpersisted gap" }), + ], + }), + }), + ]); + }); }); diff --git a/extensions/telegram/src/bot.test.ts b/extensions/telegram/src/bot.test.ts index 8b80fe4f5a1e..e89eee8df355 100644 --- a/extensions/telegram/src/bot.test.ts +++ b/extensions/telegram/src/bot.test.ts @@ -2107,6 +2107,7 @@ describe("createTelegramBot", () => { () => "telegram:default:42", ); telegramBotDepsForTest.readAmbientTranscriptWatermark = vi.fn(() => ({ + sessionId: "session-current", messageId: "502", timestampMs: 1_736_380_860_000, updatedAt: 1_736_380_900_000, diff --git a/src/auto-reply/reply/get-reply-run.media-only.test.ts b/src/auto-reply/reply/get-reply-run.media-only.test.ts index d1b12f27c6de..27f808893634 100644 --- a/src/auto-reply/reply/get-reply-run.media-only.test.ts +++ b/src/auto-reply/reply/get-reply-run.media-only.test.ts @@ -2055,6 +2055,7 @@ describe("runPreparedReply media-only handling", () => { key: '["telegram","","-100123",""]', messageId: "35676", timestampMs: 1_710_000_000_000, + expectedSessionId: expect.any(String), }); expect(call?.followupRun.currentInboundContext?.text).toContain( "#35675 obviyus ->#35674: Are you fr fr", diff --git a/src/auto-reply/reply/get-reply-run.ts b/src/auto-reply/reply/get-reply-run.ts index e0c09e5d31c7..3aa19f4e2044 100644 --- a/src/auto-reply/reply/get-reply-run.ts +++ b/src/auto-reply/reply/get-reply-run.ts @@ -174,6 +174,7 @@ function normalizeMessageTimestampMs(value: unknown): number | undefined { } async function updateRoomEventAmbientTranscriptWatermark(params: { + expectedSessionId: string; sessionCtx: TemplateContext; storePath?: string; sessionKey?: string; @@ -191,6 +192,7 @@ async function updateRoomEventAmbientTranscriptWatermark(params: { key, messageId, timestampMs: params.sessionCtx.AmbientTranscriptTimestampMs, + expectedSessionId: params.expectedSessionId, }); } @@ -1307,6 +1309,7 @@ export async function runPreparedReply( onMessagePersisted: isRoomEvent ? async () => await updateRoomEventAmbientTranscriptWatermark({ + expectedSessionId: preparedSessionState.sessionId, sessionCtx, storePath, sessionKey: sessionKey ?? preparedSessionState.sessionId, diff --git a/src/config/sessions/ambient-transcript-watermark.test.ts b/src/config/sessions/ambient-transcript-watermark.test.ts new file mode 100644 index 000000000000..85d0ee0c0612 --- /dev/null +++ b/src/config/sessions/ambient-transcript-watermark.test.ts @@ -0,0 +1,124 @@ +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { afterEach, beforeEach, describe, expect, it } from "vitest"; +import { + readAmbientTranscriptWatermark, + resolveAmbientTranscriptWatermarkKey, + updateAmbientTranscriptWatermark, +} from "./ambient-transcript-watermark.js"; +import { loadSessionEntry, replaceSessionEntry } from "./session-accessor.js"; + +describe("ambient transcript watermark", () => { + let tempDir: string; + let storePath: string; + const sessionKey = "agent:main:telegram:group:-100123"; + const key = resolveAmbientTranscriptWatermarkKey({ + channel: "telegram", + accountId: "default", + conversationId: "-100123", + }); + + beforeEach(() => { + tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "openclaw-ambient-watermark-")); + storePath = path.join(tempDir, "sessions.json"); + }); + + afterEach(() => { + fs.rmSync(tempDir, { recursive: true, force: true }); + }); + + it("stamps and resolves the watermark for the current session id only", async () => { + await replaceSessionEntry( + { sessionKey, storePath }, + { sessionId: "before-reset", updatedAt: 1_700_000_000_000 }, + ); + + await updateAmbientTranscriptWatermark({ + storePath, + sessionKey, + key, + messageId: "11", + timestampMs: 1_700_000_001_000, + }); + + const persistedEntry = loadSessionEntry({ sessionKey, storePath }); + if (!persistedEntry) { + throw new Error("Expected persisted session entry"); + } + expect(persistedEntry?.ambientTranscriptWatermarks?.[key]).toMatchObject({ + sessionId: "before-reset", + messageId: "11", + timestampMs: 1_700_000_001_000, + }); + expect(readAmbientTranscriptWatermark(persistedEntry, key)).toMatchObject({ + sessionId: "before-reset", + messageId: "11", + }); + + await replaceSessionEntry( + { sessionKey, storePath }, + { + ...persistedEntry, + sessionId: "after-reset", + updatedAt: 1_700_000_002_000, + }, + ); + + const resetEntry = loadSessionEntry({ sessionKey, storePath }); + expect(readAmbientTranscriptWatermark(resetEntry, key)).toBeUndefined(); + + await updateAmbientTranscriptWatermark({ + storePath, + sessionKey, + key, + messageId: "12", + timestampMs: 1_700_000_002_000, + expectedSessionId: "before-reset", + }); + + expect( + readAmbientTranscriptWatermark(loadSessionEntry({ sessionKey, storePath }), key), + ).toBeUndefined(); + + await updateAmbientTranscriptWatermark({ + storePath, + sessionKey, + key, + messageId: "12", + timestampMs: 1_700_000_002_000, + expectedSessionId: "after-reset", + }); + + expect( + readAmbientTranscriptWatermark(loadSessionEntry({ sessionKey, storePath }), key), + ).toMatchObject({ + sessionId: "after-reset", + messageId: "12", + }); + }); + + it("ignores legacy watermarks without a session id", () => { + fs.writeFileSync( + storePath, + JSON.stringify({ + [sessionKey]: { + sessionId: "current-session", + updatedAt: 1_700_000_000_000, + ambientTranscriptWatermarks: { + [key]: { + messageId: "11", + timestampMs: 1_700_000_001_000, + updatedAt: 1_700_000_002_000, + }, + }, + }, + }), + "utf-8", + ); + + expect( + readAmbientTranscriptWatermark(loadSessionEntry({ sessionKey, storePath }), key), + ).toBeUndefined(); + }); +}); diff --git a/src/config/sessions/ambient-transcript-watermark.ts b/src/config/sessions/ambient-transcript-watermark.ts index da18f2141600..6548f8b94a2d 100644 --- a/src/config/sessions/ambient-transcript-watermark.ts +++ b/src/config/sessions/ambient-transcript-watermark.ts @@ -52,10 +52,14 @@ function isAmbientTranscriptWatermarkAfter( } export function readAmbientTranscriptWatermark( - entry: Pick | undefined, + entry: Pick | undefined, key: string, ): AmbientTranscriptWatermark | undefined { - return entry?.ambientTranscriptWatermarks?.[key]; + const watermark = entry?.ambientTranscriptWatermarks?.[key]; + // A watermark only vouches for rows in the transcript it was written against. + // After a session reset those rows live in an archived file the model never + // reads, so a cross-session (or legacy sessionId-less) watermark must not hide them. + return watermark?.sessionId === entry?.sessionId ? watermark : undefined; } export async function updateAmbientTranscriptWatermark(params: { @@ -64,6 +68,7 @@ export async function updateAmbientTranscriptWatermark(params: { key: string; messageId: string; timestampMs?: number; + expectedSessionId?: string; }): Promise { return await updateSessionEntry( { @@ -71,6 +76,15 @@ export async function updateAmbientTranscriptWatermark(params: { sessionKey: params.sessionKey, }, (entry) => { + // onMessagePersisted fires after the durable row write; if the session was + // reset in between, stamping the new sessionId would hide rows that only + // exist in the archived transcript. Skip the advance instead. + if (!entry.sessionId) { + return null; + } + if (params.expectedSessionId !== undefined && entry.sessionId !== params.expectedSessionId) { + return null; + } const current = readAmbientTranscriptWatermark(entry, params.key); if ( !isAmbientTranscriptWatermarkAfter( @@ -84,6 +98,7 @@ export async function updateAmbientTranscriptWatermark(params: { ambientTranscriptWatermarks: { ...entry.ambientTranscriptWatermarks, [params.key]: { + sessionId: entry.sessionId, messageId: params.messageId, ...(params.timestampMs !== undefined ? { timestampMs: params.timestampMs } : {}), updatedAt: Date.now(), diff --git a/src/config/sessions/types.ts b/src/config/sessions/types.ts index 73a590c7d30a..157eabd72e10 100644 --- a/src/config/sessions/types.ts +++ b/src/config/sessions/types.ts @@ -113,6 +113,7 @@ export type SessionContextBudgetStatus = { }; export type AmbientTranscriptWatermark = { + sessionId: string; messageId: string; timestampMs?: number; updatedAt: number;