Skip to content

Commit c8dba85

Browse files
fix: keep Telegram outbound context on transcript time (#98769)
1 parent b2c507c commit c8dba85

8 files changed

Lines changed: 225 additions & 18 deletions

extensions/telegram/src/bot-message-dispatch.test.ts

Lines changed: 57 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1621,6 +1621,38 @@ describe("dispatchTelegramMessage draft streaming", () => {
16211621
expectDraftStreamParams({ maxChars: 800 });
16221622
});
16231623

1624+
it("marks durable non-preview finals with the transcript prompt-context timestamp", async () => {
1625+
const transcriptTimestamp = Date.now() + 1_000;
1626+
const context = createContext();
1627+
context.ctxPayload.SessionKey = "agent:default:telegram:direct:123";
1628+
mockDefaultSessionEntry();
1629+
readLatestAssistantTextByIdentity.mockResolvedValue({
1630+
text: "Final answer",
1631+
timestamp: transcriptTimestamp,
1632+
});
1633+
deliverInboundReplyWithMessageSendContext.mockResolvedValue({
1634+
status: "handled_visible",
1635+
delivery: {
1636+
messageIds: ["2001"],
1637+
visibleReplySent: true,
1638+
},
1639+
});
1640+
dispatchReplyWithBufferedBlockDispatcher.mockImplementation(async ({ dispatcherOptions }) => {
1641+
await dispatcherOptions.deliver({ text: "Final answer" }, { kind: "final" });
1642+
return { queuedFinal: true };
1643+
});
1644+
1645+
await dispatchWithContext({ context, streamMode: "off" });
1646+
1647+
const outbound = expectRecordFields(mockCallArg(deliverInboundReplyWithMessageSendContext), {
1648+
payload: expect.objectContaining({ text: "Final answer" }),
1649+
});
1650+
expectRecordFields(expectRecordFields(outbound.payload, {}).channelData, {
1651+
telegram: { promptContextTimestampMs: transcriptTimestamp },
1652+
});
1653+
expect(deliverReplies).not.toHaveBeenCalled();
1654+
});
1655+
16241656
it("keeps the Telegram edit cap for non-block previews regardless of chunk config", async () => {
16251657
const draftStream = createDraftStream();
16261658
createTelegramDraftStream.mockReturnValue(draftStream);
@@ -1644,12 +1676,20 @@ describe("dispatchTelegramMessage draft streaming", () => {
16441676

16451677
it("streams text-only finals into the answer message", async () => {
16461678
const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 });
1679+
const transcriptTimestamp = Date.now() + 1_000;
1680+
const context = createContext();
1681+
context.ctxPayload.SessionKey = "agent:default:telegram:direct:123";
1682+
mockDefaultSessionEntry();
1683+
readLatestAssistantTextByIdentity.mockResolvedValue({
1684+
text: "Final answer",
1685+
timestamp: transcriptTimestamp,
1686+
});
16471687
dispatchReplyWithBufferedBlockDispatcher.mockImplementation(async ({ dispatcherOptions }) => {
16481688
await dispatcherOptions.deliver({ text: "Final answer" }, { kind: "final" });
16491689
return { queuedFinal: true };
16501690
});
16511691

1652-
await dispatchWithContext({ context: createContext() });
1692+
await dispatchWithContext({ context });
16531693

16541694
expect(answerDraftStream.update).toHaveBeenCalledWith("Final answer");
16551695
expect(answerDraftStream.stop).toHaveBeenCalled();
@@ -1664,11 +1704,20 @@ describe("dispatchTelegramMessage draft streaming", () => {
16641704
messageId: 2001,
16651705
text: "Final answer",
16661706
messageThreadId: 777,
1707+
promptContextTimestampMs: transcriptTimestamp,
16671708
});
16681709
});
16691710

16701711
it("records streamed final replies into the prompt context cache", async () => {
16711712
const storePath = `/tmp/openclaw-telegram-stream-context-${process.pid}-${Date.now()}.json`;
1713+
const transcriptTimestamp = Date.now() + 1_000;
1714+
const context = createContext();
1715+
context.ctxPayload.SessionKey = "agent:default:telegram:direct:123";
1716+
mockDefaultSessionEntry();
1717+
readLatestAssistantTextByIdentity.mockResolvedValue({
1718+
text: "Done already: timeoutSeconds is now 7200s.",
1719+
timestamp: transcriptTimestamp,
1720+
});
16721721
setupDraftStreams({ answerMessageId: 1497 });
16731722
dispatchReplyWithBufferedBlockDispatcher.mockImplementation(async ({ dispatcherOptions }) => {
16741723
await dispatcherOptions.deliver(
@@ -1679,7 +1728,7 @@ describe("dispatchTelegramMessage draft streaming", () => {
16791728
});
16801729

16811730
await dispatchWithContext({
1682-
context: createContext(),
1731+
context,
16831732
cfg: { session: { store: storePath } },
16841733
telegramDeps: {
16851734
...telegramDepsForTest,
@@ -1704,7 +1753,7 @@ describe("dispatchTelegramMessage draft streaming", () => {
17041753
},
17051754
});
17061755

1707-
const context = await buildTelegramConversationContext({
1756+
const conversationContext = await buildTelegramConversationContext({
17081757
cache,
17091758
accountId: "default",
17101759
chatId: "123",
@@ -1715,10 +1764,13 @@ describe("dispatchTelegramMessage draft streaming", () => {
17151764
replyTargetWindowSize: 2,
17161765
});
17171766

1718-
expect(context.map((entry) => entry.node.messageId)).toContain("1497");
1719-
expect(context.map((entry) => entry.node.body)).toContain(
1767+
expect(conversationContext.map((entry) => entry.node.messageId)).toContain("1497");
1768+
expect(conversationContext.map((entry) => entry.node.body)).toContain(
17201769
"Done already: timeoutSeconds is now 7200s.",
17211770
);
1771+
expect(
1772+
conversationContext.find((entry) => entry.node.messageId === "1497")?.node.timestamp,
1773+
).toBe(transcriptTimestamp);
17221774
});
17231775

17241776
it("suppresses text-only tool payloads delivered after the final answer", async () => {

extensions/telegram/src/bot-message-dispatch.ts

Lines changed: 50 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -117,7 +117,10 @@ import {
117117
type LaneName,
118118
} from "./lane-delivery.js";
119119
import { TELEGRAM_TEXT_CHUNK_LIMIT } from "./outbound-adapter.js";
120-
import { recordOutboundMessageForPromptContext } from "./outbound-message-context.js";
120+
import {
121+
recordOutboundMessageForPromptContext,
122+
withTelegramPromptContextTimestampMs,
123+
} from "./outbound-message-context.js";
121124
import {
122125
createTelegramReasoningStepState,
123126
splitTelegramReasoningText,
@@ -244,6 +247,7 @@ export type TelegramDispatchResult =
244247
type TelegramReasoningLevel = "off" | "on" | "stream";
245248

246249
type TelegramTranscriptMirrorPayload = { text?: string; mediaUrls?: string[] };
250+
type CurrentTurnTranscriptFinal = { text: string; timestamp: number };
247251
type TelegramScopedTranscriptSession = { sessionId: string; storePath: string };
248252
type FreshTelegramSessionEntryLoader = ((
249253
agentId: string,
@@ -1459,10 +1463,16 @@ export const dispatchTelegramMessage = async ({
14591463
const sessionKey = ctxPayload.SessionKey;
14601464
let transcriptMirrorSequence = 0;
14611465
const transcriptMirrorTurnId = `${chatId}:${ctxPayload.MessageSid ?? msg.message_id ?? dispatchStartedAt}`;
1462-
const resolveCurrentTurnTranscriptFinalText = async (): Promise<string | undefined> => {
1466+
let currentTurnTranscriptFinal: CurrentTurnTranscriptFinal | undefined;
1467+
const resolveCurrentTurnTranscriptFinal = async (): Promise<
1468+
CurrentTurnTranscriptFinal | undefined
1469+
> => {
14631470
if (!sessionKey) {
14641471
return undefined;
14651472
}
1473+
if (currentTurnTranscriptFinal) {
1474+
return currentTurnTranscriptFinal;
1475+
}
14661476
try {
14671477
const { entry: sessionEntry, storePath } = loadFreshSessionEntry(route.agentId, sessionKey);
14681478
if (!sessionEntry?.sessionId) {
@@ -1477,12 +1487,25 @@ export const dispatchTelegramMessage = async ({
14771487
if (!latest?.timestamp || latest.timestamp < dispatchStartedAt) {
14781488
return undefined;
14791489
}
1480-
return latest.text;
1490+
currentTurnTranscriptFinal = {
1491+
text: latest.text,
1492+
timestamp: latest.timestamp,
1493+
};
1494+
return currentTurnTranscriptFinal;
14811495
} catch (err) {
14821496
logVerbose(`telegram transcript final candidate lookup failed: ${formatErrorMessage(err)}`);
14831497
return undefined;
14841498
}
14851499
};
1500+
const resolveCurrentTurnTranscriptFinalText = async (): Promise<string | undefined> =>
1501+
(await resolveCurrentTurnTranscriptFinal())?.text;
1502+
const resolvePromptContextTimestampMs = async (text: string): Promise<number | undefined> => {
1503+
const final = await resolveCurrentTurnTranscriptFinal();
1504+
if (final?.text.trim() !== text.trim()) {
1505+
return undefined;
1506+
}
1507+
return final.timestamp;
1508+
};
14861509
const deliveryBaseOptions = {
14871510
chatId: String(chatId),
14881511
accountId: route.accountId,
@@ -1631,6 +1654,14 @@ export const dispatchTelegramMessage = async ({
16311654
return false;
16321655
}
16331656
const deliverablePayload = applyQuoteReplyTarget(payload);
1657+
const promptContextTimestampMs =
1658+
options?.durable && deliverablePayload.text
1659+
? await resolvePromptContextTimestampMs(deliverablePayload.text)
1660+
: undefined;
1661+
const effectivePayload = withTelegramPromptContextTimestampMs(
1662+
deliverablePayload,
1663+
promptContextTimestampMs,
1664+
);
16341665
const silent = options?.silent ?? (silentErrorReplies && payload.isError === true);
16351666
const durableDelivery = telegramDeps.deliverInboundReplyWithMessageSendContext;
16361667
if (options?.durable && durableDelivery) {
@@ -1641,7 +1672,7 @@ export const dispatchTelegramMessage = async ({
16411672
accountId: route.accountId,
16421673
agentId: route.agentId,
16431674
ctxPayload,
1644-
payload: deliverablePayload,
1675+
payload: effectivePayload,
16451676
info: { kind: "final" },
16461677
replyToMode,
16471678
threadId: threadSpec.id,
@@ -1652,13 +1683,13 @@ export const dispatchTelegramMessage = async ({
16521683
},
16531684
silent,
16541685
requiredCapabilities: deriveDurableFinalDeliveryRequirements({
1655-
payload: deliverablePayload,
1656-
replyToId: deliverablePayload.replyToId,
1686+
payload: effectivePayload,
1687+
replyToId: effectivePayload.replyToId,
16571688
threadId: threadSpec.id,
16581689
silent,
16591690
payloadTransport: true,
16601691
extraCapabilities: {
1661-
nativeQuote: usesNativeTelegramQuote(deliverablePayload),
1692+
nativeQuote: usesNativeTelegramQuote(effectivePayload),
16621693
},
16631694
}),
16641695
});
@@ -1676,7 +1707,7 @@ export const dispatchTelegramMessage = async ({
16761707
const result = await (telegramDeps.deliverReplies ?? deliverReplies)({
16771708
...deliveryBaseOptions,
16781709
transcriptMirror: options?.durable ? deliveryBaseOptions.transcriptMirror : undefined,
1679-
replies: [deliverablePayload],
1710+
replies: [effectivePayload],
16801711
onVoiceRecording: sendRecordVoice,
16811712
silent,
16821713
mediaLoader: telegramDeps.loadWebMedia,
@@ -1701,6 +1732,10 @@ export const dispatchTelegramMessage = async ({
17011732
groupId: deliveryBaseOptions.mirrorGroupId,
17021733
});
17031734
try {
1735+
const promptContextContent =
1736+
result.delivery.promptContextContent ?? result.delivery.content;
1737+
const promptContextTimestampMs =
1738+
await resolvePromptContextTimestampMs(promptContextContent);
17041739
await (
17051740
telegramDeps.recordOutboundMessageForPromptContext ??
17061741
recordOutboundMessageForPromptContext
@@ -1710,7 +1745,8 @@ export const dispatchTelegramMessage = async ({
17101745
chatId: deliveryBaseOptions.chatId,
17111746
message: { message_id: result.delivery.messageId },
17121747
messageId: result.delivery.messageId,
1713-
text: result.delivery.promptContextContent ?? result.delivery.content,
1748+
text: promptContextContent,
1749+
...(promptContextTimestampMs !== undefined ? { promptContextTimestampMs } : {}),
17141750
...(threadSpec.id !== undefined ? { messageThreadId: threadSpec.id } : {}),
17151751
});
17161752
} catch (error) {
@@ -1868,11 +1904,13 @@ export const dispatchTelegramMessage = async ({
18681904
markProgressFinalDelivered();
18691905
return { kind: "sent" };
18701906
};
1871-
const resolveTranscriptBackedFinalText = async (text: string): Promise<string> =>
1872-
await resolveTranscriptBackedChannelFinalText({
1907+
const resolveTranscriptBackedFinalText = async (text: string): Promise<string> => {
1908+
const candidate = await resolveCurrentTurnTranscriptFinal();
1909+
return await resolveTranscriptBackedChannelFinalText({
18731910
finalText: text,
1874-
resolveCandidateText: resolveCurrentTurnTranscriptFinalText,
1911+
resolveCandidateText: async () => candidate?.text,
18751912
});
1913+
};
18761914

18771915
if (isDmTopic) {
18781916
try {

extensions/telegram/src/message-cache.ts

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -66,6 +66,9 @@ export type TelegramMessageCache = {
6666
};
6767

6868
type MessageWithExternalReply = Message & { external_reply?: Message };
69+
type MessageWithPromptContextTimestamp = Message & {
70+
openclaw_prompt_context_timestamp_ms?: unknown;
71+
};
6972

7073
type TelegramMessageCacheBucket = {
7174
messages: Map<string, TelegramCachedMessageNode>;
@@ -159,6 +162,15 @@ function resolveMediaType(placeholder?: string): string | undefined {
159162
return placeholder?.match(/^<media:([^>]+)>$/)?.[1];
160163
}
161164

165+
function resolveMessageTimestamp(msg: Message): number | undefined {
166+
const promptContextTimestamp = (msg as MessageWithPromptContextTimestamp)
167+
.openclaw_prompt_context_timestamp_ms;
168+
if (typeof promptContextTimestamp === "number" && Number.isFinite(promptContextTimestamp)) {
169+
return promptContextTimestamp;
170+
}
171+
return msg.date ? msg.date * 1000 : undefined;
172+
}
173+
162174
function normalizeMessageNode(
163175
msg: Message,
164176
params: { threadId?: number },
@@ -172,13 +184,14 @@ function normalizeMessageNode(
172184
const replyMessage = resolveReplyMessage(msg);
173185
const body = resolveMessageBody(msg);
174186
const threadId = normalizeTelegramCacheThreadId(params.threadId);
187+
const timestamp = resolveMessageTimestamp(msg);
175188
return {
176189
sourceMessage: msg,
177190
messageId: String(msg.message_id),
178191
sender: buildSenderName(msg) ?? "unknown sender",
179192
...(msg.from?.id != null ? { senderId: String(msg.from.id) } : {}),
180193
...(msg.from?.username ? { senderUsername: msg.from.username } : {}),
181-
...(msg.date ? { timestamp: msg.date * 1000 } : {}),
194+
...(timestamp !== undefined ? { timestamp } : {}),
182195
...(body ? { body } : {}),
183196
...(media ? { mediaType: resolveMediaType(media.placeholder) ?? media.placeholder } : {}),
184197
...(fileId ? { mediaRef: `telegram:file/${fileId}` } : {}),

extensions/telegram/src/outbound-adapter.test.ts

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -158,6 +158,29 @@ describe("telegramOutbound", () => {
158158
expect(result).toEqual({ channel: "telegram", messageId: "tg-buttons", chatId: "12345" });
159159
});
160160

161+
it("forwards prompt-context timestamps on durable payload sends", async () => {
162+
sendMessageTelegramMock.mockResolvedValueOnce({ messageId: "tg-final", chatId: "12345" });
163+
164+
const result = await telegramOutbound.sendPayload!({
165+
cfg: {} as never,
166+
to: "12345",
167+
text: "",
168+
payload: {
169+
text: "Final answer",
170+
channelData: {
171+
telegram: {
172+
promptContextTimestampMs: 1_779_394_740_123,
173+
},
174+
},
175+
},
176+
deps: { sendTelegram: sendMessageTelegramMock },
177+
});
178+
179+
const options = callOptionsAt(sendMessageTelegramMock, 0, "12345", "Final answer");
180+
expect(options.promptContextTimestampMs).toBe(1_779_394_740_123);
181+
expect(result).toEqual({ channel: "telegram", messageId: "tg-final", chatId: "12345" });
182+
});
183+
161184
it("applies reaction-only payloads without sending empty Telegram text", async () => {
162185
reactMessageTelegramMock.mockResolvedValueOnce({ ok: true });
163186

extensions/telegram/src/outbound-adapter.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import type { TelegramInlineButtons } from "./button-types.js";
2626
import { resolveTelegramInlineButtons } from "./button-types.js";
2727
import { splitTelegramHtmlChunks } from "./format.js";
2828
import { resolveTelegramInteractiveTextFallback } from "./interactive-fallback.js";
29+
import { resolveTelegramPromptContextTimestampMs } from "./outbound-message-context.js";
2930
import { parseTelegramReplyToMessageId, parseTelegramThreadId } from "./outbound-params.js";
3031
import { loadTelegramSendModule, type TelegramSendModule } from "./send-runtime.js";
3132
import { normalizeTelegramOutboundTarget, parseTelegramTarget } from "./targets.js";
@@ -156,6 +157,7 @@ export async function sendTelegramPayloadMessages(params: {
156157
const payloadOpts = {
157158
...params.baseOpts,
158159
quoteText,
160+
promptContextTimestampMs: resolveTelegramPromptContextTimestampMs(params.payload),
159161
...(params.payload.audioAsVoice === true ? { asVoice: true } : {}),
160162
};
161163
const shouldConsumeImplicitReplyTarget =

0 commit comments

Comments
 (0)