diff --git a/extensions/telegram/src/bot-message-dispatch.test.ts b/extensions/telegram/src/bot-message-dispatch.test.ts index 44b6339741c0..3adfcb3813c4 100644 --- a/extensions/telegram/src/bot-message-dispatch.test.ts +++ b/extensions/telegram/src/bot-message-dispatch.test.ts @@ -411,6 +411,18 @@ describe("dispatchTelegramMessage draft streaming", () => { return expectRecordFields(mockCallArg(dispatchReplyWithBufferedBlockDispatcher), expected); } + // The collapse bar edits the live window message in place (finalizeToPreview) + // instead of deleting it and reposting the bar as a new message. + function expectWindowCollapsedTo( + stream: { finalizeToPreview: { mock: { calls: unknown[][] } } }, + barText: string, + ) { + const calls = stream.finalizeToPreview.mock.calls; + expect(calls.length).toBeGreaterThan(0); + const preview = calls[calls.length - 1][0] as { text?: string }; + expect(preview.text).toBe(barText); + } + function createContext(overrides?: Partial): TelegramMessageContext { const base = { ctxPayload: {}, @@ -2888,11 +2900,12 @@ describe("dispatchTelegramMessage draft streaming", () => { ); expect(answerDraftStream.update).not.toHaveBeenCalledWith("Branch is up to date"); expect(answerDraftStream.forceNewMessage).toHaveBeenCalledTimes(1); - expect(answerDraftStream.clear).toHaveBeenCalledTimes(1); - // The progress window collapses to a one-line activity summary (Discord - // parity) before the final answer posts fresh below it. - expectDeliveredReply(0, { text: "🛠️ 1 tool call · ⏱️ 1s" }); - expectDeliveredReply(0, { text: "Branch is up to date" }, 1); + // The window collapses IN PLACE into the one-line activity summary (edit, + // not delete + repost — Discord parity), so clear() is never called on it. + expect(answerDraftStream.clear).not.toHaveBeenCalled(); + expectWindowCollapsedTo(answerDraftStream, "🛠️ 1 tool call · ⏱️ 1s"); + // The final answer then posts fresh below the collapsed bar. + expectDeliveredReply(0, { text: "Branch is up to date" }); expect(editMessageTelegram).not.toHaveBeenCalled(); }); @@ -2905,7 +2918,7 @@ describe("dispatchTelegramMessage draft streaming", () => { } it("tallies reasoning bursts and tool calls into the collapse summary", async () => { - setupDraftStreams({ answerMessageId: 2001 }); + const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 }); dispatchReplyWithBufferedBlockDispatcher.mockImplementation( async ({ dispatcherOptions, replyOptions }) => { // burst 1 → tool → burst 2 → tool, then a trailing burst flushed at the @@ -2928,8 +2941,8 @@ describe("dispatchTelegramMessage draft streaming", () => { telegramCfg: { streaming: { mode: "progress" } }, }); - expectDeliveredReply(0, { text: "🧠 3 thoughts · 🛠️ 2 tool calls · ⏱️ 1s" }); - expectDeliveredReply(0, { text: "Done" }, 1); + expectWindowCollapsedTo(answerDraftStream, "🧠 3 thoughts · 🛠️ 2 tool calls · ⏱️ 1s"); + expectDeliveredReply(0, { text: "Done" }); }); it("does not post a collapse summary when no progress draft started", async () => { @@ -2974,6 +2987,88 @@ describe("dispatchTelegramMessage draft streaming", () => { expect(texts.some((text) => text.includes("tool call · ⏱️"))).toBe(false); }); + it("keeps the progress window alive under /reasoning on so commentary and tools still stream", async () => { + // /reasoning on removes only the 🧠 lane from the window; commentary, tool + // lines, and the collapse bar must still stream (Discord parity). A prior + // regression forced block streaming in progress mode, killing the window. + loadSessionStore.mockReturnValue({ s1: { reasoningLevel: "on" } }); + const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 }); + dispatchReplyWithBufferedBlockDispatcher.mockImplementation( + async ({ dispatcherOptions, replyOptions }) => { + await replyOptions?.onItemEvent?.({ kind: "preamble", itemId: "c1", progressText: "Note" }); + await replyOptions?.onToolStart?.({ name: "exec", phase: "start" }); + await dispatcherOptions.deliver({ text: "Done" }, { kind: "final" }); + return { queuedFinal: true }; + }, + ); + + await dispatchWithContext({ + context: createContext({ + ctxPayload: { SessionKey: "s1" } as unknown as TelegramMessageContext["ctxPayload"], + }), + streamMode: "progress", + telegramCfg: { streaming: { mode: "progress" } }, + }); + + // The window streamed (a preview was rendered) and collapsed into a bar + // counting the note + tool — proof the window was not killed. + expect(answerDraftStream.updatePreview).toHaveBeenCalled(); + expectWindowCollapsedTo(answerDraftStream, "💬 1 note · 🛠️ 1 tool call · ⏱️ 1s"); + expectDeliveredReply(0, { text: "Done" }); + }); + + it("does not duplicate tool lines into the window under verbose", async () => { + // Invariant D2 (persistent XOR window): when the durable verbose lane owns + // tool messages, the window must render no tool line and must not count it. + const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 }); + dispatchReplyWithBufferedBlockDispatcher.mockImplementation( + async ({ dispatcherOptions, replyOptions }) => { + replyOptions?.onVerboseProgressVisibility?.(true); + await replyOptions?.onToolStart?.({ name: "exec", phase: "start" }); + await dispatcherOptions.deliver({ text: "Done" }, { kind: "final" }); + return { queuedFinal: true }; + }, + ); + + await dispatchWithContext({ + context: createContext(), + streamMode: "progress", + telegramCfg: { streaming: { mode: "progress" } }, + }); + + // No tool line ever rendered to the window (verbose owns it durably), so the + // window never streamed and there is no collapse bar to count it. + expect(answerDraftStream.updatePreview).not.toHaveBeenCalled(); + expect(answerDraftStream.finalizeToPreview).not.toHaveBeenCalled(); + const texts = allDeliveredReplyTexts(); + expect(texts.some((text) => text.includes("tool call"))).toBe(false); + }); + + it("posts a collapse summary for a message_tool_only final that bypasses the answer path", async () => { + // Codex-runtime turns deliver the final out-of-band (queuedFinal), so the + // in-band collapse path never runs. The window still started, so the + // cleanup-time fallback must emit the bar (Discord parity). + setupDraftStreams({ answerMessageId: 2001 }); + dispatchReplyWithBufferedBlockDispatcher.mockImplementation(async ({ replyOptions }) => { + await replyOptions?.onItemEvent?.({ kind: "preamble", itemId: "c1", progressText: "Note" }); + await replyOptions?.onToolStart?.({ name: "exec", phase: "start" }); + return { + queuedFinal: true, + counts: { block: 0, final: 1, tool: 1 }, + sourceReplyDeliveryMode: "message_tool_only", + }; + }); + + await dispatchWithContext({ + context: createContext(), + streamMode: "progress", + telegramCfg: { streaming: { mode: "progress" } }, + }); + + const texts = allDeliveredReplyTexts(); + expect(texts).toContain("💬 1 note · 🛠️ 1 tool call · ⏱️ 1s"); + }); + it("replaces Telegram command progress items with matching command output", async () => { const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 }); dispatchReplyWithBufferedBlockDispatcher.mockImplementation(async ({ replyOptions }) => { @@ -3037,9 +3132,10 @@ describe("dispatchTelegramMessage draft streaming", () => { expect(answerDraftStream.forceNewMessage.mock.invocationCallOrder[1]).toBeLessThan( answerDraftStream.update.mock.invocationCallOrder[0], ); - // Collapse summary posts first, then the final answer below it. - expectDeliveredReply(0, { text: "🛠️ 1 tool call · ⏱️ 1s" }); - expectDeliveredReply(0, { text: "Branch is up to date" }, 1); + // Window collapses in place into the summary bar; the final answer posts + // fresh below it. + expectWindowCollapsedTo(answerDraftStream, "🛠️ 1 tool call · ⏱️ 1s"); + expectDeliveredReply(0, { text: "Branch is up to date" }); }); it("does not stream text-only tool results into progress drafts", async () => { @@ -3122,8 +3218,8 @@ describe("dispatchTelegramMessage draft streaming", () => { expect(answerDraftStream.updatePreview).toHaveBeenCalledWith( telegramProgressPreview("Shelling\n\n🛠️ Exec", "Shelling\n🛠️ Exec"), ); - expectDeliveredReply(0, { text: "🛠️ 1 tool call · ⏱️ 1s" }); - expectDeliveredReply(0, { text: "Branch is up to date" }, 1); + expectWindowCollapsedTo(answerDraftStream, "🛠️ 1 tool call · ⏱️ 1s"); + expectDeliveredReply(0, { text: "Branch is up to date" }); }); it("does not restart progress drafts for command output after final answer delivery", async () => { @@ -3153,8 +3249,8 @@ describe("dispatchTelegramMessage draft streaming", () => { expect(answerDraftStream.updatePreview).toHaveBeenCalledWith( telegramProgressPreview("Shelling\n\n🛠️ Exec", "Shelling\n🛠️ Exec"), ); - expectDeliveredReply(0, { text: "🛠️ 1 tool call · ⏱️ 1s" }); - expectDeliveredReply(0, { text: "Branch is up to date" }, 1); + expectWindowCollapsedTo(answerDraftStream, "🛠️ 1 tool call · ⏱️ 1s"); + expectDeliveredReply(0, { text: "Branch is up to date" }); }); it("does not restart progress drafts for command output while final answer delivery is pending", async () => { @@ -3188,12 +3284,12 @@ describe("dispatchTelegramMessage draft streaming", () => { expect(answerDraftStream.updatePreview).toHaveBeenCalledWith( telegramProgressPreview("Shelling\n\n🛠️ Exec", "Shelling\n🛠️ Exec"), ); - expectDeliveredReply(0, { text: "🛠️ 1 tool call · ⏱️ 1s" }); - expectDeliveredReply(0, { text: "Branch is up to date" }, 1); + expectWindowCollapsedTo(answerDraftStream, "🛠️ 1 tool call · ⏱️ 1s"); + expectDeliveredReply(0, { text: "Branch is up to date" }); }); it("uses the transcript final when progress-mode final text is truncated", async () => { - setupDraftStreams({ answerMessageId: 2001 }); + const { answerDraftStream } = setupDraftStreams({ answerMessageId: 2001 }); const fullAnswer = "Ja. Hier nochmal sauber Schritt fuer Schritt. Einen API Key kopiert man aus der Google Cloud Console. Danach pruefst du die Projekt- und API-Einstellungen."; const truncatedFinal = @@ -3219,8 +3315,8 @@ describe("dispatchTelegramMessage draft streaming", () => { telegramCfg: { streaming: { mode: "progress" } }, }); - expectDeliveredReply(0, { text: "🛠️ 1 tool call · ⏱️ 1s" }); - expectDeliveredReply(0, { text: fullAnswer }, 1); + expectWindowCollapsedTo(answerDraftStream, "🛠️ 1 tool call · ⏱️ 1s"); + expectDeliveredReply(0, { text: fullAnswer }); }); it("streams the first long final chunk and sends follow-up chunks", async () => { @@ -3949,7 +4045,7 @@ describe("dispatchTelegramMessage draft streaming", () => { await dispatchWithContext({ context: createReasoningStreamContext() }); - expect(reasoningDraftStream.update).toHaveBeenCalledWith("Thinking\n\n_Thinking_"); + expect(reasoningDraftStream.update).toHaveBeenCalledWith("🧠 _Thinking_"); expect(answerDraftStream.update).toHaveBeenCalledWith("Answer"); expect(deliverReplies).not.toHaveBeenCalled(); }); @@ -3969,7 +4065,7 @@ describe("dispatchTelegramMessage draft streaming", () => { await dispatchWithContext({ context: createReasoningForumTopicContext() }); - expect(reasoningDraftStream.update).toHaveBeenCalledWith("Thinking\n\n_Thinking_"); + expect(reasoningDraftStream.update).toHaveBeenCalledWith("🧠 _Thinking_"); expect(answerDraftStream.update).toHaveBeenCalledWith("Answer"); expect(answerDraftStream.stop).toHaveBeenCalled(); expect(deliverReplies).not.toHaveBeenCalled(); @@ -4018,7 +4114,7 @@ describe("dispatchTelegramMessage draft streaming", () => { await dispatchWithContext({ context: createReasoningStreamContext() }); expect(reasoningDraftStream.update).toHaveBeenLastCalledWith( - "Thinking\n\n_Reading_\n\n_Checking_", + "🧠 _Reading_\n\n_Checking_", ); const updates = reasoningDraftStream.update.mock.calls.map((call) => call[0]); expect(updates.join("\n")).not.toContain("CheckingReading"); @@ -4047,7 +4143,7 @@ describe("dispatchTelegramMessage draft streaming", () => { }, }); - expect(reasoningDraftStream.update).toHaveBeenCalledWith("Thinking\n\n_Thinking_"); + expect(reasoningDraftStream.update).toHaveBeenCalledWith("🧠 _Thinking_"); expect(answerDraftStream.update).toHaveBeenCalledWith("Answer"); }); @@ -4068,10 +4164,12 @@ describe("dispatchTelegramMessage draft streaming", () => { const run = dispatchWithContext({ context: createReasoningStreamContext() }); await vi.waitFor(() => - expect(reasoningDraftStream.update).toHaveBeenCalledWith("Thinking\n\n_Thinking_"), + expect(reasoningDraftStream.update).toHaveBeenCalledWith("🧠 _Thinking_"), ); + // Durable thoughts render behind the 🧠 marker; the literal "Thinking" + // header (and its streaming dot-variants) must never leak back into a lane. + expect(reasoningDraftStream.update).not.toHaveBeenCalledWith("Thinking\n\n_Thinking_"); expect(reasoningDraftStream.update).not.toHaveBeenCalledWith("Thinking.\n\n_Thinking_"); - expect(reasoningDraftStream.update).not.toHaveBeenCalledWith("Thinking..\n\n_Thinking_"); expect(reasoningDraftStream.update).not.toHaveBeenCalledWith("Thinking...\n\n_Thinking_"); finishRun?.(); await run; @@ -4158,7 +4256,7 @@ describe("dispatchTelegramMessage draft streaming", () => { await dispatchWithContext({ context: createReasoningStreamContext() }); - expect(reasoningDraftStream.update).toHaveBeenCalledWith("Thinking\n\n_hidden_"); + expect(reasoningDraftStream.update).toHaveBeenCalledWith("🧠 _hidden_"); expect(deliverReplies).not.toHaveBeenCalled(); }); @@ -4180,7 +4278,7 @@ describe("dispatchTelegramMessage draft streaming", () => { }), }); - const delivered = expectDeliveredReply(0, { text: "Thinking\n\n_hidden_" }); + const delivered = expectDeliveredReply(0, { text: "🧠 _hidden_" }); expect(delivered).not.toHaveProperty("isReasoning"); }); diff --git a/extensions/telegram/src/bot-message-dispatch.ts b/extensions/telegram/src/bot-message-dispatch.ts index ebd8b9cdb812..f720b3c23811 100644 --- a/extensions/telegram/src/bot-message-dispatch.ts +++ b/extensions/telegram/src/bot-message-dispatch.ts @@ -900,7 +900,14 @@ export const dispatchTelegramMessage = async ({ agentId: route.agentId, loadFreshSessionEntry, }); - const forceBlockStreamingForReasoning = resolvedReasoningLevel === "on"; + // Progress mode's ephemeral working-lane window IS the streaming mechanism and + // is independent of reasoning persistence (Discord keeps its window alive + // regardless of /reasoning). Only non-progress modes upgrade reasoning-on to + // block streaming. Forcing block streaming in progress mode killed the whole + // window (no commentary/tool lanes, no collapse bar) and suppressed all + // streamed output for message_tool_only providers. + const forceBlockStreamingForReasoning = + resolvedReasoningLevel === "on" && streamMode !== "progress"; const streamReasoningDraft = resolvedReasoningLevel === "stream"; const streamDeliveryEnabled = !isRoomEvent && streamMode !== "off"; const rawReplyQuoteText = @@ -1087,13 +1094,15 @@ export const dispatchTelegramMessage = async ({ }); let finalAnswerDeliveryStarted = false; let finalAnswerDelivered = false; - // While the durable verbose lane is active, the ephemeral draft yields its - // commentary lines so they render once. Tool/plan status lines keep the - // draft: they have no durable counterpart in streamed runs. + // While the durable verbose lane is active it owns EVERY progress surface + // (commentary, tool, plan, command output, patch summaries), posting each as + // its own persistent message. The ephemeral window must therefore render none + // of them, or each renders twice (invariant: persistent message XOR window). let verboseProgressActive: () => boolean = () => false; const canPushStreamToolProgress = () => Boolean( answerLane.stream && + !verboseProgressActive() && !answerLane.finalized && !finalAnswerDeliveryStarted && !finalAnswerDelivered, @@ -1130,6 +1139,7 @@ export const dispatchTelegramMessage = async ({ }; const markProgressFinalDelivered = () => { finalAnswerDelivered = true; + sawProgressFinal = true; progressDraft.markFinalReplyDelivered(); }; const resetProgressDraftState = () => { @@ -1561,6 +1571,11 @@ export const dispatchTelegramMessage = async ({ const silentErrorReplies = telegramCfg.silentErrorReplies === true; const isDmTopic = !isGroup && threadSpec.scope === "dm" && threadSpec.id != null; let queuedFinal = false; + // A final answer was produced this turn (in-band or out-of-band). Out-of-band + // finals (message_tool_only / codex) never flow through + // deliverProgressModeFinalAnswer, so the collapse bar must be posted from the + // cleanup fallback instead — see the finally block. + let sawProgressFinal = false; let skippedDuplicateAnswerBlockDraftDelivery = false; let suppressSilentReplyFallback = false; let hadErrorReplyFailureOrSkip = false; @@ -1901,44 +1916,82 @@ export const dispatchTelegramMessage = async ({ await emitPreviewFinalizedHook(result); return result.kind !== "skipped"; }; - // Post-turn collapse summary (Discord parity): when the progress window - // collapses at end-of-turn, post a one-line activity digest as a durable - // standalone message, then the final answer posts below it so the timeline - // reads thoughts/tools → summary → answer. Emitted at most once per turn, - // only for a non-error final, and only when the window actually rendered - // (rv mode delivers everything durably and the window stays empty — no bar). - const deliverProgressCollapseSummary = async () => { + // The one-line activity digest for the collapse bar, or undefined when the + // window never rendered (rv mode delivers everything durably — no bar) or + // the summary was already emitted this turn. + const resolveProgressCollapseSummaryLine = (): string | undefined => { if (progressSummaryDelivered) { - return; + return undefined; } progressSummaryDelivered = true; if (!progressDraftEverRendered) { - return; + return undefined; } const line = formatTelegramProgressSummaryLine( progressSummary.counts(), Date.now() - progressSummaryStartedAt, ); + return line || undefined; + }; + // Post-turn collapse summary (Discord parity) as a durable standalone + // message. Used when there is no live window to collapse in place — the + // final answer posts below so the timeline reads thoughts/tools → summary → + // answer. Emitted at most once per turn. + const deliverProgressCollapseSummary = async () => { + const line = resolveProgressCollapseSummaryLine(); if (!line) { return; } await sendPayload({ text: line }, { durable: true }); }; + // Collapse the live window IN PLACE into the summary bar: edit the existing + // window message so its content becomes the bar line, keeping it on screen. + // Mirrors Discord — deleting the window and reposting the bar scroll-jumps + // the Telegram client and flashes the window away. Returns true when the + // window was collapsed in place; false when there is no bar (nothing + // streamed) or no live window message, so the caller tears the window down. + const collapseProgressWindowIntoSummary = async (): Promise => { + const line = resolveProgressCollapseSummaryLine(); + if (!line) { + return false; + } + const messageId = await answerLane.stream?.finalizeToPreview(renderStreamText(line)); + if (typeof messageId === "number") { + return true; + } + // No live window to edit (rv mode, never rendered): keep the bar as a + // fresh durable post so the timeline still shows the collapse summary. + await sendPayload({ text: line }, { durable: true }); + return false; + }; const deliverProgressModeFinalAnswer = async ( payload: ReplyPayload, text: string, ): Promise => { - if (activeAnswerDraftIsToolProgressOnly) { - await rotateAnswerLaneAfterToolProgress(); - } else { - await answerLane.stream?.clear(); - resetDraftLaneState(answerLane); - } + // Collapse the window into the bar in place BEFORE resetting lane state + // (which drops the stream's message id). Error finals get no summary + // (Discord parity). When nothing collapsed in place, tear the window down + // so a stale progress box does not linger above the final answer. + const collapsedInPlace = + payload.isError === true ? false : await collapseProgressWindowIntoSummary(); if (payload.isError === true) { - // Error finals get no collapse summary (Discord parity); mark it handled. progressSummaryDelivered = true; + } + if (!collapsedInPlace) { + if (activeAnswerDraftIsToolProgressOnly) { + await rotateAnswerLaneAfterToolProgress(); + } else { + await answerLane.stream?.clear(); + resetDraftLaneState(answerLane); + } } else { - await deliverProgressCollapseSummary(); + if (activeAnswerDraftIsToolProgressOnly) { + resetAnswerToolProgressDraft(); + suppressProgressDraftState(); + rotateAnswerLaneWhenQueuedBlocksSettle = false; + } + answerLane.stream?.forceNewMessage(); + resetDraftLaneState(answerLane); } const delivered = await sendPayload(applyTextToPayload(payload, text), { durable: true }); if (!delivered) { @@ -2521,7 +2574,10 @@ export const dispatchTelegramMessage = async ({ // own durable messages and must NOT also feed the bar (invariant: // persistent message XOR bar count — D2). if (payload.phase === "start") { - if (verboseProgressActive() || !canPushStreamToolProgress()) { + // canPushStreamToolProgress() is false under verbose, so this + // also closes bursts (never counting the tool) when the durable + // lane owns the tool message (invariant: persistent XOR window). + if (!canPushStreamToolProgress()) { progressSummary.closeReasoningBurst(); progressSummary.closeCommentaryBurst(); } else { @@ -2682,6 +2738,12 @@ export const dispatchTelegramMessage = async ({ return { kind: "completed" }; } ({ queuedFinal } = turnResult.dispatchResult); + // Out-of-band finals (message_tool_only) never run the in-band final-delivery + // path, so record the final from the dispatch counts for the cleanup-time + // collapse-bar fallback. + if ((turnResult.dispatchResult.counts?.final ?? 0) > 0) { + sawProgressFinal = true; + } suppressSilentReplyFallback = turnResult.dispatchResult.sourceReplyDeliveryMode === "message_tool_only"; } catch (err) { @@ -2710,6 +2772,20 @@ export const dispatchTelegramMessage = async ({ await stream.clear(); } } + // Fallback collapse summary (Discord parity): finals that bypass + // deliverProgressModeFinalAnswer — notably message_tool_only/codex turns + // whose final is delivered out-of-band — still collapse here. The internal + // once-guard and progressDraftEverRendered check keep this from + // double-posting or firing when the window never rendered. + if ( + streamMode === "progress" && + sawProgressFinal && + !dispatchError && + !hadErrorReplyFailureOrSkip && + !isDispatchSuperseded() + ) { + await deliverProgressCollapseSummary(); + } } } finally { dispatchWasSuperseded = isDispatchSuperseded(); diff --git a/extensions/telegram/src/draft-stream.test-helpers.ts b/extensions/telegram/src/draft-stream.test-helpers.ts index eb92fe22b502..d60a755eaad9 100644 --- a/extensions/telegram/src/draft-stream.test-helpers.ts +++ b/extensions/telegram/src/draft-stream.test-helpers.ts @@ -14,6 +14,9 @@ type TestDraftStream = { stop: ReturnType Promise>>; discard: ReturnType Promise>>; materialize: ReturnType Promise>>; + finalizeToPreview: ReturnType< + typeof vi.fn<(preview: TelegramDraftPreview) => Promise> + >; forceNewMessage: ReturnType void>>; sendMayHaveLanded: ReturnType boolean>>; setMessageId: (value: number | undefined) => void; @@ -66,6 +69,15 @@ export function createTestDraftStream(params?: { await params?.onDiscard?.(); }), materialize: vi.fn().mockImplementation(async () => messageId), + finalizeToPreview: vi.fn().mockImplementation(async (preview: TelegramDraftPreview) => { + if (messageId == null) { + return undefined; + } + previewRevision += 1; + lastDeliveredText = preview.text.trimEnd(); + stopped = true; + return messageId; + }), forceNewMessage: vi.fn().mockImplementation(() => { stopped = false; if (params?.clearMessageIdOnForceNew) { @@ -113,6 +125,14 @@ export function createSequencedTestDraftStream(startMessageId = 1001): TestDraft stop: vi.fn().mockResolvedValue(undefined), discard: vi.fn().mockResolvedValue(undefined), materialize: vi.fn().mockImplementation(async () => activeMessageId), + finalizeToPreview: vi.fn().mockImplementation(async (preview: TelegramDraftPreview) => { + if (activeMessageId == null) { + return undefined; + } + previewRevision += 1; + lastDeliveredText = preview.text.trimEnd(); + return activeMessageId; + }), forceNewMessage: vi.fn().mockImplementation(() => { activeMessageId = undefined; visibleSinceMs = undefined; diff --git a/extensions/telegram/src/draft-stream.ts b/extensions/telegram/src/draft-stream.ts index 78ed1d9d7208..351905a6bcd4 100644 --- a/extensions/telegram/src/draft-stream.ts +++ b/extensions/telegram/src/draft-stream.ts @@ -59,6 +59,13 @@ export type TelegramDraftStream = { discard?: () => Promise; /** Return the current preview message id after pending updates settle. */ materialize?: () => Promise; + /** + * Collapse the preview in place: edit the existing window message so its + * content becomes `preview`, then stop without deleting. Used at end-of-turn + * so the streaming window becomes the summary bar (no delete + repost, which + * scroll-jumps the client). Returns the message id if the edit landed. + */ + finalizeToPreview: (preview: TelegramDraftPreview) => Promise; /** Reset internal state so the next update creates a new message instead of editing. */ forceNewMessage: () => void; /** True when a preview sendMessage was attempted but the response was lost. */ @@ -587,6 +594,28 @@ export function createTelegramDraftStream(params: { return streamMessageId; }; + const finalizeToPreview = async ( + preview: TelegramDraftPreview, + ): Promise => { + // Settle pending updates so we edit the real, current window message. + streamState.final = true; + await loop.flush(); + const text = preview.text.trimEnd(); + // No live window message to edit (never rendered, or already torn down): + // nothing to collapse in place — caller falls back to a fresh bar post. + if (typeof streamMessageId !== "number" || !text) { + return undefined; + } + // Replace the whole message with the bar line: edits diff from a zero + // offset, not from the streamed prefix. + deliveredTextOffset = 0; + lastSentPreviewKey = ""; + lastRequestedText = text; + lastRequestedPreview = { ...preview, text }; + await sendOrEditStreamMessage(text); + return streamMessageId; + }; + params.log?.(`telegram stream preview ready (maxChars=${maxChars}, throttleMs=${throttleMs})`); return { @@ -601,6 +630,7 @@ export function createTelegramDraftStream(params: { stop, discard, materialize, + finalizeToPreview, forceNewMessage, sendMayHaveLanded: () => messageSendAttempted && typeof streamMessageId !== "number", }; diff --git a/extensions/telegram/src/reasoning-lane-coordinator.ts b/extensions/telegram/src/reasoning-lane-coordinator.ts index 6d2d12fc5658..6d68440e833e 100644 --- a/extensions/telegram/src/reasoning-lane-coordinator.ts +++ b/extensions/telegram/src/reasoning-lane-coordinator.ts @@ -5,8 +5,23 @@ import { normalizeLowercaseStringOrEmpty } from "openclaw/plugin-sdk/string-coer import { findCodeRegions, isInsideCode } from "openclaw/plugin-sdk/text-chunking"; import { stripReasoningTagsFromText } from "openclaw/plugin-sdk/text-chunking"; -const REASONING_MESSAGE_RE = /^Thinking\.{0,3}\s*_/u; +// A durable reasoning message already marked channel-side: 🧠 + italic body +// (see markReasoningMessage). Detect it so a re-split passes it through +// unchanged instead of re-marking. +const REASONING_MESSAGE_RE = /^🧠\s+_/u; +// Core's formatReasoningMessage prefixes the italic body with a literal +// "Thinking" header. Telegram renders durable thoughts with the 🧠 marker +// (Discord parity), so this header must be rewritten channel-side. +const CORE_THINKING_HEADER_RE = /^Thinking\.{0,3}\s*\n+/u; const LEGACY_REASONING_MESSAGE_PREFIX = "Reasoning:\n"; + +// Rewrite core's "Thinking\n\n_body_" into "🧠 _body_": strip the header word +// and prefix the first italic line with 🧠. Keeps the italic body intact so +// Telegram HTML renders it as before. +function markReasoningMessage(formatted: string): string { + const withoutHeader = formatted.replace(CORE_THINKING_HEADER_RE, ""); + return withoutHeader.replace(/^_/u, "🧠 _"); +} const REASONING_TAG_PREFIXES = [ " LEGACY_REASONING_MESSAGE_PREFIX.length @@ -94,7 +114,11 @@ export function splitTelegramReasoningText( const taggedReasoning = extractThinkingFromTaggedStreamOutsideCode(text); const strippedAnswer = stripReasoningTagsFromText(text, { mode: "strict", trim: "both" }); - return { reasoningText: formatReasoningMessage(taggedReasoning || strippedAnswer || text) }; + return { + reasoningText: markReasoningMessage( + formatReasoningMessage(taggedReasoning || strippedAnswer || text), + ), + }; } type BufferedFinalAnswer = {