/** * ResponseStreamer β€” renders a whole agent turn into as FEW Telegram messages as * possible, edited at most once per throttle window (anti-spam, avoids 429s). * * The turn is modelled as ordered segments so the transcript reads clearly: * β€’ plain prose = the agent talking to you * β€’ > πŸ’­ quoted block = the agent's thinking * β€’ πŸ”§ + code block = tool calls / terminal commands / diffs * * A single "live" message is edited as content grows; only when it would exceed * Telegram's size limit is it sealed and a new live message started. */ import type { Api } from "grammy"; import { chunkMarkdown } from "../render/chunk.js"; import { toTelegramMarkdown } from "../render/markdown.js"; import { extractProgress, progressBar } from "../render/progress.js"; import { estimateProgress } from "../render/progress-estimate.js"; import { stripTelegramActionFences } from "../render/telegram-bridge.js"; import { truncateMiddle } from "../render/truncate.js"; import { safeEdit, safeSend } from "../bot/telegram-io.js"; import { outboundThreadExtra } from "../forum/thread.js"; const SOFT_LIMIT = 3500; /** Display budget for a thinking block (middle-truncated; session context keeps all). */ const THINK_DISPLAY_MAX = 2800; type SegKind = "out" | "think" | "tool"; interface Seg { kind: SegKind; text: string; /** When set, later tool updates replace this segment instead of appending. */ toolId?: string; } export interface StreamerOptions { /** Chat-like mode: drop thoughts/tools/plan; only stream agent prose. */ proseOnly?: boolean; /** When false, never render a progress bar (manager chat). Default true. */ showProgressBar?: boolean; /** * Pre-posted message id to edit in place (e.g. General "Thinking…" placeholder). * Avoids a separate bubble when the first real tokens arrive. */ seedMessageId?: number; } export class ResponseStreamer { private readonly segs: Seg[] = []; private sealedIdx = 0; private liveId: number | undefined; private timer: NodeJS.Timeout | undefined; private dirty = false; private flushing = false; private closed = false; /** Latest task-progress % parsed from the agent's `{progress: N%}` markers * (sticky across flushes; rendered as a bar on the live message). */ private progress: number | undefined; /** True once the agent emitted a real `{progress}` marker β€” from then on its * values are authoritative and the bot fallback stops contributing. */ private agentReported = false; /** Real work signals for the fallback estimate (monotonic within a turn). */ private toolCalls = 0; private outChars = 0; private thoughtChars = 0; /** * Active plan board (ACP sessionUpdate "plan"). Always rendered just above * the progress bar when set β€” done / in-progress / pending steps. */ private planMarkdown: string | undefined; private readonly proseOnly: boolean; private readonly showProgressBar: boolean; constructor( private readonly api: Api, private readonly chatId: number, private readonly throttleMs: number, private replyTo?: number, private footer?: string, private readonly onProgress?: (pct: number) => void, /** Show a bot-computed bar when the agent emits no marker. */ private readonly fallbackEnabled = false, /** Turn start time, used by the fallback's elapsed-time signal. */ private readonly turnStartedAt = Date.now(), /** Forum topic thread β€” required so stream edits land in the right topic. */ private readonly messageThreadId?: number, opts?: StreamerOptions, ) { this.proseOnly = !!opts?.proseOnly; this.showProgressBar = opts?.showProgressBar !== false; if (opts?.seedMessageId !== undefined) this.liveId = opts.seedMessageId; } /** Replace the hashtag footer (used after a logical fork swaps the session id * mid-turn, so the streamed response carries the NEW session's tags). */ setFooter(footer: string): void { this.footer = footer; } /** Seed/replace the live bubble id (General Thinking… placeholder). */ seedLiveMessage(messageId: number): void { this.liveId = messageId; } /** Current live Telegram message id (for attaching suggestions after finalize). */ get liveMessageId(): number | undefined { return this.liveId; } /** "\n\n