diff --git a/docs/guide.md b/docs/guide.md index 5d08613..4157717 100644 --- a/docs/guide.md +++ b/docs/guide.md @@ -8,7 +8,7 @@ Telegram in real time. One paired DM owner controls the bridge; optional group chat access remains separately configured from the terminal, never by the model. - **Inbound:** DMs / group mentions → injected as `` user turns (photos attached inline; other files downloaded to an inbox). -- **Outbound:** assistant output streams live — native message **drafts** for DMs (Bot API 9.3+), **edited-message** previews for groups — then one finalized MarkdownV2 message per turn. A headless host can turn all of that off with `set profile daemon`, leaving `telegram_send` / `telegram_ask` as the only way out. +- **Outbound:** assistant output streams live through native message drafts for DMs (Bot API 9.3+) and edited-message previews for groups. Each turn ends with a MarkdownV2 message by default, or optional rich Markdown on Bot API 10.1+. A headless host can turn automatic output off with `set profile daemon`, leaving `telegram_send` / `telegram_ask` as the only way out. - **Control:** local `/telegram` configuration, owner-only Telegram commands (`/spawn`, `/sessions`, `/cleanup`, `/stop`, `/status`), and three model tools (`telegram_send`, `telegram_react`, `telegram_ask`). - **Zero runtime dependencies** — the raw Bot API over Bun's `fetch`/`FormData`. @@ -185,10 +185,11 @@ configured groups never receive process-spawning authority. | Key | Values | Default | |---|---|---| | `streaming` | `true` (live preview) \| `false` (per turn) \| `final` (last message only) \| `explicit` (nothing automatic) | `true` | +| `richMessages` | `off` (MarkdownV2) \| `auto` (rich constructs) \| `on` (prefer rich Markdown) | `off` | | `profile` | `daemon` (headless host: forces `explicit`, no idle notify post, `telegram_ask` always on) \| `default` | `default` | | `deliverAs` | `steer` \| `followUp` — how inbound queues while the agent is busy | `followUp` | | `chunkMode` | `length` \| `newline` | `newline` | -| `textChunkLimit` | `1`–`4096` | `4096` | +| `textChunkLimit` | `1`–`4096`; also caps rich source when explicitly set | unset (4096 legacy; 32768 whole rich) | | `replyToMode` | `off` \| `first` \| `all` — threading for `telegram_send` replies | `first` | | `ackReaction` | a whitelist emoji (empty to disable) — reaction on receipt | unset | | `mentionPatterns` | JSON array of regexes that also satisfy group mention-gating, e.g. `["\\bassistant\\b"]` | unset | @@ -261,8 +262,9 @@ count as a mention. ## Model tools - **`telegram_send`** — send text and/or files to the active or durably claimed - chat (or a given `chat_id`). Text is chunked and rendered as MarkdownV2 - (plain-text fallback on parse errors). `files` are absolute paths: images send + chat (or a given `chat_id`). Markdown follows `richMessages`, with MarkdownV2 + as the default and plain-text fallback on parse errors. `format: "text"` stays + literal in every mode. `files` are absolute paths: images send as photos, everything else as documents (≤ 50 MB each). - **`telegram_react`** — react to a message with a Telegram whitelist emoji (👍 👎 ❤ 🔥 👀 🎉 …). @@ -370,12 +372,28 @@ transcript. task's closing line and not an answer to anyone. A missing reply is visible to the person waiting and can be asked again; a leaked internal turn cannot be recalled. Ask for an answer, not a transcript. -- Any reply too long for one Telegram message (4096 chars, or `textChunkLimit`) - is split at a paragraph, line, or word boundary (`chunkMode`) and each part is - prefixed `(i/n)`, so it arrives complete and in order rather than cut off. - Code fences are closed and reopened across the split. This covers assistant - answers and command output alike — a long `/sessions` listing or `/cleanup` - preview splits too, with the keyboard on the final part. +- Legacy messages split at 4096 characters (or `textChunkLimit`), with room + reserved for formatting and labels. Splits prefer paragraph, line, or word + boundaries (`chunkMode`). Each part carries `(i/n)`; code fences close and + reopen across splits. Command output keeps this limit, with keyboards on the + final part. +- `set richMessages auto` selects rich Markdown for tables, task lists, + `
`, paired `$$` math, and `` outside code. Headings alone + don't select it. `on` prefers rich Markdown for all Markdown output; `off` + keeps MarkdownV2. Rich messages require Bot API 10.1+. +- Rich delivery sends the original Markdown, including Telegram's native task + syntax. Telegram controls rendering and checkbox interaction. This setting + adds no checklist-management commands or state store. +- Live drafts and edit previews stay unchanged. Final messages and final + preview edits use the selected format. An answer with no permanent preview or + committed prefix can arrive as one rich message up to 32768 UTF-16 source + units when `textChunkLimit` is unset. An explicit cap remains authoritative. + Longer answers and existing previews keep legacy chunk boundaries. +- A definitive rich rejection (400 or unsupported-method 404) falls back to + MarkdownV2, then plain text on a parse rejection. A rejected whole answer is + split again to fit legacy limits. Network failures, timeouts, server errors, + authorization failures, and exhausted rate limits don't trigger another-format + send. Constructs split across parts may lose formatting. - A part Telegram rate-limits (`429`) is retried up to three times, honouring `retry_after`, instead of dropping the rest of the answer. diff --git a/src/access.ts b/src/access.ts index fce078f..fb2bced 100644 --- a/src/access.ts +++ b/src/access.ts @@ -37,7 +37,7 @@ export type Access = { ackReaction?: string; /** Which chunks carry Telegram's reply reference. Default "first". */ replyToMode?: "off" | "first" | "all"; - /** Max chars per outbound message before splitting. Default 4096, clamp 1..4096. */ + /** Explicit source cap, clamped to 1..4096. Unset: 4096 legacy, 32768 whole rich. */ textChunkLimit?: number; /** Split strategy. Default "newline". */ chunkMode?: "length" | "newline"; @@ -60,6 +60,8 @@ export type Access = { * recalled. */ streaming?: boolean | "final" | "explicit"; + /** Final Markdown formatting. Missing = off (MarkdownV2). */ + richMessages?: "auto" | "on" | "off"; /** * Headless-host contract, set via `/telegram set profile daemon`. Absent = * `default` (an interactive laptop session). @@ -92,9 +94,8 @@ export type Access = { }; /** - * Per-message character budget for this config, clamped to Telegram's cap. - * Every outbound text path splits against this, so a long reply is never - * rejected whole or silently cut. + * Legacy per-message budget, also used for previews and rich fallback chunks. + * Whole rich messages may use a larger budget only when no explicit cap is set. */ export function messageLimit(access: Access): number { return Math.max(1, Math.min(access.textChunkLimit ?? TELEGRAM_MAX_CHARS, TELEGRAM_MAX_CHARS)); @@ -272,6 +273,10 @@ export function loadAccess(warn?: (msg: string) => void): Access { ? parsed.streaming : undefined; const profile: Access["profile"] = parsed.profile === "daemon" ? "daemon" : undefined; + const richMessages: Access["richMessages"] = + parsed.richMessages === "auto" || parsed.richMessages === "on" || parsed.richMessages === "off" + ? parsed.richMessages + : undefined; return { enabled: parsed.enabled ?? false, dmPolicy: parsed.dmPolicy ?? "pairing", @@ -285,6 +290,7 @@ export function loadAccess(warn?: (msg: string) => void): Access { chunkMode: parsed.chunkMode, deliverAs: parsed.deliverAs, streaming, + richMessages, profile, transcribeCommand: Array.isArray(parsed.transcribeCommand) && parsed.transcribeCommand.every((arg) => typeof arg === "string") ? parsed.transcribeCommand diff --git a/src/index.ts b/src/index.ts index 662a25e..42e9cbf 100644 --- a/src/index.ts +++ b/src/index.ts @@ -154,6 +154,7 @@ const TELEGRAM_ARGS: CompletionNode = { mentionPatterns: null, deliverAs: { steer: null, followUp: null }, streaming: { true: null, false: null, final: null, explicit: null }, + richMessages: { auto: null, on: null, off: null }, profile: { daemon: null, default: null }, transcribeCommand: null, }, @@ -192,6 +193,7 @@ const SET_KEY_HELP: Record = { mentionPatterns: "JSON array of mention regexes", deliverAs: "steer | followUp delivery", streaming: "output: true (stream) | false (per-turn) | final (one message) | explicit (tool calls only)", + richMessages: "formatting: auto (rich constructs) | on (prefer rich) | off (MarkdownV2)", profile: "daemon (headless: explicit output + always-on telegram_ask) | default", transcribeCommand: "JSON argv for voice transcription", }; @@ -1493,6 +1495,7 @@ export default function telegramExtension(pi: ExtensionAPI): void { // the streaming value is overridden, and "off" would have been read as // "silent" when it only ever meant "no live preview". `Streaming: ${STREAMING_LABEL[String(effectiveStreaming(a))]} · profile: ${a.profile ?? "default"} · deliverAs: ${a.deliverAs ?? "followUp"} · chunkMode: ${a.chunkMode ?? "newline"} · replyTo: ${a.replyToMode ?? "first"}`, + `Rich messages: ${a.richMessages ?? "off"}`, `Notify: ${a.notifyMode ?? "off"}${a.notifyChat ? ` · chat ${a.notifyChat}` : ""}`, `Voice transcription: ${a.transcribeCommand?.length ? a.transcribeCommand.join(" ") : "off"}`, `Control topic: ${a.controlThreadId != null ? `#${a.controlThreadId}` : "not attached"}`, @@ -1834,6 +1837,9 @@ export default function telegramExtension(pi: ExtensionAPI): void { return ctx.ui.notify("streaming: true | false | final | explicit", "warning"); } a.streaming = value === "final" || value === "explicit" ? value : value === "true"; + } else if (key === "richMessages") { + if (value !== "auto" && value !== "on" && value !== "off") return ctx.ui.notify("richMessages: auto | on | off", "warning"); + a.richMessages = value; } else if (key === "profile") { if (value !== "daemon" && value !== "default") return ctx.ui.notify("profile: daemon | default", "warning"); a.profile = value === "daemon" ? "daemon" : undefined; @@ -1850,7 +1856,7 @@ export default function telegramExtension(pi: ExtensionAPI): void { } } } else { - return ctx.ui.notify(`set: unknown key "${key}". Keys: ackReaction, replyToMode, textChunkLimit, chunkMode, mentionPatterns, deliverAs, streaming, profile, transcribeCommand`, "warning"); + return ctx.ui.notify(`set: unknown key "${key}". Keys: ackReaction, replyToMode, textChunkLimit, chunkMode, mentionPatterns, deliverAs, streaming, richMessages, profile, transcribeCommand`, "warning"); } saveAccess(a); access = a; diff --git a/src/index.wiring.test.ts b/src/index.wiring.test.ts index 77eb8e9..61762df 100644 --- a/src/index.wiring.test.ts +++ b/src/index.wiring.test.ts @@ -117,6 +117,61 @@ async function startBridge(h: Harness): Promise { } describe("extension wiring", () => { + test("rich formatting survives command reload, rejects invalid input, and respects literal/file sends", async () => { + writeAccess({ enabled: true, allowFrom: ["42"], profile: "daemon" }); + const first = harness(["read"]); + const ctx = { ui: { notify() {} } }; + await first.commands.get("telegram")!.handler("set richMessages auto", ctx); + await first.commands.get("telegram")!.handler("set richMessages AUTO", ctx); + expect(loadAccess().richMessages).toBe("auto"); + + const h = harness(["read"]); + const calls: Array<{ method: string; payload?: Record }> = []; + const filesDir = mkdtempSync(join(tmpdir(), "omp-tg-rich-files-")); + const previousFetch = globalThis.fetch; + globalThis.fetch = (async (input, init) => { + const method = String(input).split("/").pop()!; + calls.push({ method, payload: typeof init?.body === "string" ? JSON.parse(init.body) : undefined }); + return new Response(JSON.stringify({ ok: true, result: { message_id: calls.length } })); + }) as typeof fetch; + try { + await startBridge(h); + const source = "| Name | State |\n| --- | --- |\n| build | ready |"; + calls.length = 0; + const tool = h.tools.get("telegram_send")!; + await tool.execute("rich", { chat_id: "42", text: source }, undefined, undefined, {}); + expect(calls).toEqual([{ method: "sendRichMessage", payload: { chat_id: "42", rich_message: { markdown: source } } }]); + await h.commands.get("telegram")!.handler("set richMessages on", ctx); + calls.length = 0; + await tool.execute("literal", { chat_id: "42", text: source, format: "text" }, undefined, undefined, {}); + expect(calls).toEqual([{ method: "sendMessage", payload: { chat_id: "42", text: source } }]); + + writeFileSync(join(dir, "access.json"), JSON.stringify({ ...loadAccess(), richMessages: "corrupt" })); + await h.commands.get("telegram")!.handler("status", ctx); // Reload hand-edited policy. + calls.length = 0; + await tool.execute("corrupt", { chat_id: "42", text: source }, undefined, undefined, {}); + expect(calls[0].method).toBe("sendMessage"); + expect(calls[0].payload?.parse_mode).toBe("MarkdownV2"); + + const file = join(filesDir, "report.txt"); + writeFileSync(file, "report"); + calls.length = 0; + const result = await tool.execute("file", { chat_id: "42", text: "", files: [file] }, undefined, undefined, {}); + expect(result.isError).toBeUndefined(); + expect(calls.map((c) => c.method)).toEqual(["sendDocument"]); + calls.length = 0; + const denied = await tool.execute("denied", { chat_id: "43", text: source }, undefined, undefined, {}); + expect(denied.isError).toBe(true); + expect(calls).toEqual([]); + } finally { + rmSync(filesDir, { recursive: true, force: true }); + await h.handlers.get("session_shutdown")?.[0]?.({ type: "session_shutdown" }, { + sessionManager: { getSessionId: () => "session-1", getSessionFile: () => "/tmp/session-1.jsonl" }, + }); + globalThis.fetch = previousFetch; + } + }); + test("registers telegram_ask and the /away command", () => { const h = harness(["ask", "read"]); expect(h.tools.has("telegram_ask")).toBe(true); diff --git a/src/markdown.test.ts b/src/markdown.test.ts index 765742f..4772505 100644 --- a/src/markdown.test.ts +++ b/src/markdown.test.ts @@ -1,5 +1,5 @@ import { test, expect, describe } from "bun:test"; -import { escapeMdV2, mdToMarkdownV2, chunk, chunkLabeled, PART_LABEL_RESERVE, TELEGRAM_MAX_CHARS, MARKDOWN_HEADROOM } from "./markdown"; +import { escapeMdV2, mdToMarkdownV2, chunk, chunkLabeled, PART_LABEL_RESERVE, TELEGRAM_MAX_CHARS, MARKDOWN_HEADROOM, hasRichConstructs } from "./markdown"; // The exact MarkdownV2 special set that escapeMdV2 must prefix with a backslash. const SPECIALS = ["_", "*", "[", "]", "(", ")", "~", "`", ">", "#", "+", "-", "=", "|", "{", "}", ".", "!", "\\"]; @@ -7,6 +7,35 @@ const SPECIALS = ["_", "*", "[", "]", "(", ")", "~", "`", ">", "#", "+", "-", "= /** Count non-overlapping triple-backtick sequences — the fence-balance measure chunk() uses. */ const countFences = (s: string): number => (s.match(/```/g) ?? []).length; +describe("hasRichConstructs", () => { + test.each([ + "| Name | State |\n| --- | :---: |\n| build | ready |", + "Name | State\n--- | ---:", + "| escaped \\| pipe | state |\n| --- | --- |", + "- [ ] Pending", "+ [x] Done", "* [X]", + "
More
", + "star", + "$$E = mc^2$$", "$$\na + b\n$$", + "````\n```\n````\n- [x] outside", + "prefix `unmatched\n- [ ] outside", + ])("detects rich syntax: %s", (source) => expect(hasRichConstructs(source)).toBe(true)); + + test.each([ + "## Heading\n**ordinary bold**", "---", "Name\n---", + "| A | B |\n| --- |", "| A |\n| -- |", + "- [x]word", "\\- [x] literal", "- \\[x] literal", + "\\
", "\\", "\\$\\$math\\$\\$", + " - [x] indented", "\t
", " \t
", + "```md\n- [x] code\n```", "~~~\n
\n~~~", + "````\n```\n- [x] still code\n````", + "~~~\n```\n\n", "```\n- [x] unclosed", + "`
`", "`` `
``", + "`multiline\n
\nspan`", "$$$$", "$$ \n $$", + "| a \\| b |\n| --- | --- |", + "`| a | b |`\n| --- | --- |", + ])("ignores literal or malformed syntax: %s", (source) => expect(hasRichConstructs(source)).toBe(false)); +}); + describe("escapeMdV2", () => { test("escapes every MarkdownV2 special character with a backslash", () => { for (const c of SPECIALS) { diff --git a/src/markdown.ts b/src/markdown.ts index 6b695ac..e22742f 100644 --- a/src/markdown.ts +++ b/src/markdown.ts @@ -9,9 +9,69 @@ /** Telegram's hard per-message character cap (UTF-16 units). */ export const TELEGRAM_MAX_CHARS = 4096; +/** Conservative raw-source budget for Telegram rich messages. */ +export const TELEGRAM_RICH_MAX_CHARS = 32768; /** Headroom to reserve when a chunk will be MarkdownV2-escaped (escaping grows text). */ export const MARKDOWN_HEADROOM = 96; +/** Detect rich-only constructs without rewriting the source sent to Telegram. */ +export function hasRichConstructs(markdown: string): boolean { + let fenceChar = ""; + let fenceLength = 0; + const visible: string[] = []; + for (const line of markdown.split("\n")) { + const fence = /^ {0,3}(`{3,}|~{3,})(.*)$/.exec(line); + if (fenceLength) { + if (fence && fence[1][0] === fenceChar && fence[1].length >= fenceLength && !fence[2].trim()) fenceLength = 0; + visible.push(""); + continue; + } + if (fence && !(fence[1][0] === "`" && fence[2].includes("`"))) { + fenceChar = fence[1][0]; + fenceLength = fence[1].length; + visible.push(""); + continue; + } + // Mask escaped punctuation, preserving cell content but not its syntax. + visible.push(/^( {4}| {0,3}\t)/.test(line) ? "" : line.replace(/\\[!-/:-@[-`{-~]/g, "\0")); + } + + const source = visible.join("\n"); + const ticks = [...source.matchAll(/`+/g)]; + const next = new Map(); + const closes: Array = new Array(ticks.length); + for (let i = ticks.length - 1; i >= 0; i--) { + closes[i] = next.get(ticks[i][0].length); + next.set(ticks[i][0].length, i); + } + const fragments: string[] = []; + let offset = 0; + for (let i = 0; i < ticks.length; i++) { + const close = closes[i]; + if (close === undefined) continue; // unmatched backticks are literal + const start = ticks[i].index; + const end = ticks[close].index + ticks[close][0].length; + fragments.push(source.slice(offset, start), source.slice(start, end).replace(/[^\n]/g, "\0")); + offset = end; + i = close; + } + fragments.push(source.slice(offset)); + const text = fragments.join(""); + if (/<(?:details|tg-emoji)(?:\s[^>]*|)>/i.test(text)) return true; + for (const match of text.matchAll(/\$\$((?:(?!\$\$)[\s\S])+)\$\$/g)) { + if (match[1].replace(/\0/g, "").trim()) return true; + } + let header: string[] | undefined; + for (const line of text.split("\n")) { + if (/^ {0,3}[-+*]\s+\[[ xX]\](?:\s|$)/.test(line)) return true; + const trimmed = line.trim(); + const cells = trimmed.includes("|") ? trimmed.replace(/^\||\|$/g, "").split("|").map((cell) => cell.trim()) : undefined; + if (header && cells && cells.length === header.length && cells.every((cell) => /^:?-{3,}:?$/.test(cell))) return true; + header = cells; + } + return false; +} + /** Escape every MarkdownV2 special character with a backslash. */ export function escapeMdV2(s: string): string { return s.replace(/[_*[\]()~`>#+\-=|{}.!\\]/g, "\\$&"); diff --git a/src/notify.test.ts b/src/notify.test.ts index 80a0fe2..21588b5 100644 --- a/src/notify.test.ts +++ b/src/notify.test.ts @@ -110,6 +110,17 @@ describe("loadAccess field preservation", () => { expect(loadAccess().notifyMode).toBeUndefined(); }); + test("reloads rich formatting modes and discards invalid saved values", () => { + for (const richMessages of ["auto", "on", "off"] as const) { + saveAccess({ ...defaultAccess(), richMessages }); + expect(loadAccess().richMessages).toBe(richMessages); + } + for (const richMessages of ["ON", true, 1, null]) { + writeFileSync(join(dir, "access.json"), JSON.stringify({ ...defaultAccess(), richMessages })); + expect(loadAccess().richMessages ?? "off").toBe("off"); + } + }); + test("round-trips the streaming mode and daemon profile", () => { saveAccess({ ...defaultAccess(), streaming: "explicit", profile: "daemon" }); const a = loadAccess(); diff --git a/src/outbound.test.ts b/src/outbound.test.ts index db63f3d..5412aec 100644 --- a/src/outbound.test.ts +++ b/src/outbound.test.ts @@ -1,10 +1,15 @@ import { afterEach, test, expect, describe, setSystemTime } from "bun:test"; -import { defaultAccess } from "./access"; +import { type Access, defaultAccess } from "./access"; import { Outbound, assistantText, finalAssistantText } from "./outbound"; +import { mdToMarkdownV2 } from "./markdown"; const assistant = (text: string): unknown => ({ role: "assistant", content: [{ type: "text", text }] }); const toolResult = (): unknown => ({ role: "toolResult", content: [{ type: "text", text: "tool output" }] }); const originalFetch = globalThis.fetch; +/** Flush pending microtasks — the fetch double resolves without real I/O. */ +const flush = async (): Promise => { + for (let i = 0; i < 200; i++) await Promise.resolve(); +}; afterEach(() => { globalThis.fetch = originalFetch; @@ -286,10 +291,6 @@ describe("Outbound long answers", () => { }; /** Recorded Telegram text back to source form: drop MarkdownV2 escapes and the (i/n) label. */ const unlabel = (text: string): string => text.replace(/\\/g, "").replace(/^\(\d+\/\d+\)\n/, ""); - /** Flush pending microtasks — the fetch double resolves without real I/O. */ - const flush = async (): Promise => { - for (let i = 0; i < 200; i++) await Promise.resolve(); - }; test("a 9k answer is delivered whole, as labelled consecutive parts", async () => { const sent: string[] = []; @@ -401,3 +402,266 @@ describe("Outbound long answers", () => { outbound.shutdown(); }); }); + +describe("Outbound rich Markdown", () => { + const table = "| Name | State |\n| --- | --- |\n| build | ready |"; + type Payload = { + chat_id?: string; + text?: string; + parse_mode?: string; + message_id?: number; + message_thread_id?: number; + reply_parameters?: { message_id: number }; + rich_message?: { markdown: string }; + draft_id?: number; + }; + type Call = { method: string; payload: Payload }; + const rejected = (code: number, description = "rejected"): Response => + new Response(JSON.stringify({ ok: false, error_code: code, description, parameters: code === 429 ? { retry_after: 1 } : undefined })); + + function wire(over: Partial = {}, reject?: (call: Call) => Response | undefined) { + const calls: Call[] = []; + const chat = new Map(); + const waits: number[] = []; + let id = 0; + globalThis.fetch = (async (url, init) => { + const call = { method: String(url).split("/").pop()!, payload: JSON.parse(String(init?.body)) as Payload }; + calls.push(call); + const error = reject?.(call); + if (error) return error; + const { method, payload } = call; + const body = payload.rich_message?.markdown ?? payload.text ?? ""; + if (method === "sendMessage" || method === "sendRichMessage") chat.set(++id, body); + if (method === "editMessageText") chat.set(payload.message_id!, body); + return new Response(JSON.stringify({ ok: true, result: { message_id: id } })); + }) as typeof fetch; + const outbound = new Outbound(() => ({ ...defaultAccess(), richMessages: "auto", ...over }), undefined, async (ms) => { waits.push(ms); }); + outbound.setToken("111:rich-test"); + return { outbound, calls, chat, waits }; + } + + test("auto preserves rich source, topic and reply while off and literal text bypass rich", async () => { + const rich = wire(); + expect(await rich.outbound.send("42", table, { replyTo: 3, threadId: 9 })).toEqual([1]); + expect(rich.calls).toEqual([{ method: "sendRichMessage", payload: { + chat_id: "42", rich_message: { markdown: table }, message_thread_id: 9, reply_parameters: { message_id: 3 }, + } }]); + for (const [richMessages, format, body] of [ + ["off", "markdown", mdToMarkdownV2(table)], + ["on", "text", table], + ] as const) { + const w = wire({ richMessages }); + await w.outbound.send("42", table, { format }); + expect(w.calls[0].method).toBe("sendMessage"); + expect([...w.chat.values()]).toEqual([body]); + expect(w.calls[0].payload.parse_mode).toBe(format === "markdown" ? "MarkdownV2" : undefined); + } + for (const text of ["**bold**", "```\n- [x] code\n```", "`
`"]) { + const w = wire(); + await w.outbound.send("42", text); + expect(w.calls[0].method).toBe("sendMessage"); + expect([...w.chat.values()]).toEqual([mdToMarkdownV2(text)]); + } + const empty = wire({ richMessages: "on" }); + expect(await empty.outbound.send("42", "")).toEqual([]); + expect(empty.calls).toEqual([]); + }); + + test.each([400, 404])("rich rejection %s falls back with original source and routing", async (code) => { + const w = wire({}, ({ method }) => method === "sendRichMessage" ? rejected(code) : undefined); + expect(await w.outbound.send("42", table, { replyTo: 3, threadId: 9 })).toEqual([1]); + expect([...w.chat.values()]).toEqual([mdToMarkdownV2(table)]); + expect(w.calls.map((c) => c.method)).toEqual(["sendRichMessage", "sendMessage"]); + expect(w.calls[1].payload).toMatchObject({ message_thread_id: 9, reply_parameters: { message_id: 3 }, parse_mode: "MarkdownV2" }); + const plain = wire({}, ({ method, payload }) => method === "sendRichMessage" || payload.parse_mode ? rejected(code === 404 && method === "sendRichMessage" ? 404 : 400) : undefined); + await plain.outbound.send("42", table, { replyTo: 3, threadId: 9 }); + expect([...plain.chat.values()]).toEqual([table]); + expect(plain.calls.at(-1)?.payload).toEqual({ chat_id: "42", text: table, message_thread_id: 9, reply_parameters: { message_id: 3 } }); + }); + + test("rich rate limits retry one delivery; ambiguous and authorization failures never resend", async () => { + let limited = false; + const w = wire({}, () => { if (!limited) { limited = true; return rejected(429); } }); + await w.outbound.send("42", table); + expect(w.waits).toEqual([1250]); + expect([...w.chat.values()]).toEqual([table]); + expect(w.calls.map((c) => c.method)).toEqual(["sendRichMessage", "sendRichMessage"]); + for (const failure of [401, 403, 500, 429, "timeout"] as const) { + const failed = wire({}, () => { + if (failure === "timeout") throw new Error("network timeout"); + return rejected(failure); + }); + await expect(failed.outbound.send("42", table)).rejects.toThrow(); + expect(failed.chat.size).toBe(0); + expect(failed.calls.every((c) => c.method === "sendRichMessage")).toBe(true); + } + }); + + test("whole rich reports use the larger budget and definitive rejection re-splits once", async () => { + const source = table + "\n" + "report ".repeat(850); + const rich = wire(); + expect(await rich.outbound.send("42", source)).toEqual([1]); + expect([...rich.chat.values()]).toEqual([source]); + for (const code of [400, 404]) { + const w = wire({ chunkMode: "length" }, ({ method, payload }) => + method === "sendRichMessage" ? rejected(code) : payload.parse_mode ? rejected(400) : undefined); + expect(await w.outbound.send("42", source, { replyTo: 3, threadId: 9 })).toEqual([1, 2]); + expect(w.calls.filter((c) => c.method === "sendRichMessage")).toHaveLength(1); + const messages = [...w.chat.values()]; + expect(messages.map((m) => m.match(/^\((\d\/\d)\)\n/)?.[1])).toEqual(["1/2", "2/2"]); + expect(messages.every((m) => m.length <= 4096)).toBe(true); + expect(messages.map((m) => m.replace(/^\(\d+\/\d+\)\n/, "")).join("")).toBe(source); + expect(w.calls.filter((c) => c.method === "sendMessage").every((c) => c.payload.text!.length <= 4096)).toBe(true); + expect(w.calls.every((c) => c.payload.message_thread_id === 9)).toBe(true); + expect(w.calls.at(-1)?.payload.reply_parameters).toBeUndefined(); + } + }); + + test("explicit caps and over-budget sources keep labeled chunks and reply modes", async () => { + for (const [textChunkLimit, count, replyToMode] of [[1000, 1400, "all"], [undefined, 17000, "off"]] as const) { + const source = "日🌱".repeat(count); + const w = wire({ richMessages: "on", textChunkLimit, chunkMode: "length", replyToMode }); + const ids = await w.outbound.send("42", source, { replyTo: 3 }); + const messages = [...w.chat.values()]; + expect(ids).toEqual([...w.chat.keys()]); + expect(messages.length).toBeGreaterThan(1); + expect(messages.every((m) => m.length <= (textChunkLimit ?? 4096))).toBe(true); + expect(messages.map((m, i) => m.startsWith(`(${i + 1}/${messages.length})\n`))).toEqual(messages.map(() => true)); + expect(messages.map((m) => m.replace(/^\(\d+\/\d+\)\n/, "")).join("")).toBe(source); + expect(w.calls.every((c) => c.payload.reply_parameters?.message_id === (replyToMode === "all" ? 3 : undefined))).toBe(true); + } + const fence = wire({ richMessages: "on", textChunkLimit: 1000 }); + await fence.outbound.send("42", "```ts\n" + "const value = 1;\n".repeat(180) + "```"); + expect([...fence.chat.values()].every((m) => (m.match(/```/g) ?? []).length % 2 === 0)).toBe(true); + expect([...fence.chat.values()].join("\n").match(/const value = 1;/g)).toHaveLength(180); + }); + + test("topic recovery is shared by whole rich and legacy fallback sends", async () => { + const w = wire({}, ({ method, payload }) => { + if (payload.message_thread_id === 9) return rejected(400, "Bad Request: message thread not found"); + if (method === "sendRichMessage") return rejected(404); + }); + const recovered: number[] = []; + w.outbound.setMissingThreadHandler(async (_chat, thread) => { recovered.push(thread); return 10; }); + expect(await w.outbound.send("42", table, { threadId: 9, replyTo: 4 })).toEqual([1]); + expect(recovered).toEqual([9]); + expect(w.calls.map((c) => [c.method, c.payload.message_thread_id])).toEqual([ + ["sendRichMessage", 9], ["sendRichMessage", 10], ["sendMessage", 10], + ]); + expect(w.calls.at(-1)?.payload.reply_parameters).toEqual({ message_id: 4 }); + const failed = wire({}, ({ method }) => rejected(400, method === "sendRichMessage" ? "rich parse error" : "message thread not found")); + let recoveries = 0; + failed.outbound.setMissingThreadHandler(async () => { recoveries++; return 10; }); + await expect(failed.outbound.send("42", table, { threadId: 9 })).rejects.toThrow(); + expect(recoveries).toBe(1); + expect(failed.calls.map((c) => c.method)).toEqual(["sendRichMessage", "sendMessage", "sendMessage"]); + }); + + test("DM drafts finalize once as rich and clear the same draft in the same topic", async () => { + const source = table + "\n" + "report ".repeat(850); + const w = wire(); + try { + w.outbound.markActive("42", 9); + w.outbound.onMessageUpdate(assistant(source)); + await w.outbound.onTurnEnd(assistant(source)); + await w.outbound.onAgentEnd(source); + expect([...w.chat.values()]).toEqual([source]); + const drafts = w.calls.filter((c) => c.method === "sendMessageDraft"); + expect(drafts).toHaveLength(2); + expect(drafts[0].payload.text).toBe(source.slice(-4096)); + expect(drafts[1].payload).toEqual({ ...drafts[0].payload, text: "" }); + expect(drafts[0].payload.message_thread_id).toBe(9); + expect(w.calls.filter((c) => c.method === "sendRichMessage")).toHaveLength(1); + } finally { w.outbound.shutdown(); } + }); + + test("overflowed group previews retain every segment once with final rich labels", async () => { + const source = table + "\n" + "report ".repeat(1400); + const w = wire({ richMessages: "on", chunkMode: "length" }); + try { + setSystemTime(new Date(1_000_000)); + w.outbound.markActive("-100", 9); + w.outbound.onMessageUpdate(assistant(source.slice(0, 500))); + await flush(); + setSystemTime(new Date(1_005_000)); + w.outbound.onMessageUpdate(assistant(source)); + await flush(); + await w.outbound.onTurnEnd(assistant(source)); + await w.outbound.onAgentEnd(); + const messages = [...w.chat.values()]; + expect(messages).toHaveLength(3); + expect(messages.map((m, i) => m.startsWith(`(${i + 1}/3)\n`))).toEqual([true, true, true]); + expect(messages.map((m) => m.replace(/^\(\d+\/\d+\)\n/, "")).join("")).toBe(source); + expect(messages.some((m) => m.includes("▍"))).toBe(false); + expect(w.calls.filter((c) => c.method === "sendMessage")).toHaveLength(1); + expect(w.calls.filter((c) => c.payload.rich_message).every((c) => !c.payload.text && !c.payload.parse_mode)).toBe(true); + } finally { w.outbound.shutdown(); } + }); + + test("rich edit rejection falls back in place and not-modified never strips formatting", async () => { + for (const error of [400, 404, "not modified"] as const) { + const w = wire({}, ({ method, payload }) => { + if (method === "editMessageText" && payload.rich_message) { + return rejected(typeof error === "number" ? error : 400, error === "not modified" ? "Bad Request: message is not modified" : "rich unsupported"); + } + }); + try { + w.outbound.markActive("-100", 9); + w.outbound.onMessageUpdate(assistant(table)); + await w.outbound.onTurnEnd(assistant(table)); + await w.outbound.onAgentEnd(); + expect(w.chat.size).toBe(1); + expect(w.calls.filter((c) => c.method === "sendMessage")).toHaveLength(1); + expect(w.calls.some((c) => c.method === "sendRichMessage")).toBe(false); + const edits = w.calls.filter((c) => c.method === "editMessageText"); + expect(edits).toHaveLength(error === "not modified" ? 1 : 2); + if (error !== "not modified") expect([...w.chat.values()]).toEqual([mdToMarkdownV2(table)]); + } finally { w.outbound.shutdown(); } + } + }); + + test("finalization recovers missing topics and keeps rejected rich fallback disabled", async () => { + for (const richRejected of [false, true]) { + const w = wire({ streaming: false }, ({ method, payload }) => { + if (method === "sendChatAction") return; + if (richRejected && method === "sendRichMessage") return rejected(404); + if (payload.message_thread_id === 9) return rejected(400, "message thread not found"); + }); + let recoveries = 0; + w.outbound.setMissingThreadHandler(async () => { recoveries++; return 10; }); + try { + w.outbound.markActive("42", 9); + await w.outbound.onTurnEnd(assistant(table)); + await w.outbound.onAgentEnd(); + expect(recoveries).toBe(1); + expect([...w.chat.values()]).toEqual([richRejected ? mdToMarkdownV2(table) : table]); + expect(w.calls.at(-1)?.payload.message_thread_id).toBe(10); + expect(w.calls.filter((c) => c.method === "sendRichMessage")).toHaveLength(richRejected ? 1 : 2); + } finally { w.outbound.shutdown(); } + } + }); + + test("explicit and daemon gates stay silent while explicit sends honor rich mode", async () => { + for (const over of [{ streaming: "explicit" }, { profile: "daemon", streaming: true }] as const) { + const w = wire(over); + try { + w.outbound.markActive("42"); + w.outbound.onMessageUpdate(assistant(table)); + await w.outbound.onTurnEnd(assistant(table)); + await w.outbound.onAgentEnd(table); + expect(w.calls.filter((c) => c.method !== "sendChatAction")).toEqual([]); + await w.outbound.send("42", table); + expect([...w.chat.values()]).toEqual([table]); + expect(w.calls.at(-1)?.method).toBe("sendRichMessage"); + } finally { w.outbound.shutdown(); } + } + const boundary = wire({ richMessages: "on" }); + try { + boundary.outbound.markActive("-100", 9); + boundary.outbound.onMessageUpdate(assistant(table)); + await boundary.outbound.onSessionBoundary(); + expect([...boundary.chat.values()]).toEqual([table]); + expect(boundary.calls.some((c) => c.payload.rich_message || c.payload.parse_mode)).toBe(false); + } finally { boundary.outbound.shutdown(); } + }); +}); diff --git a/src/outbound.ts b/src/outbound.ts index 37f519b..a8e2b93 100644 --- a/src/outbound.ts +++ b/src/outbound.ts @@ -12,7 +12,7 @@ import { stat } from "node:fs/promises"; import { extname } from "node:path"; import { type Access, assertSendable, effectiveStreaming, messageLimit } from "./access"; import { isMissingThreadError, type Logger, TgError, tg, tgUpload, withRateLimit } from "./api"; -import { MARKDOWN_HEADROOM, PART_LABEL_RESERVE, TELEGRAM_MAX_CHARS, chunkLabeled, mdToMarkdownV2 } from "./markdown"; +import { MARKDOWN_HEADROOM, PART_LABEL_RESERVE, TELEGRAM_MAX_CHARS, TELEGRAM_RICH_MAX_CHARS, chunkLabeled, hasRichConstructs, mdToMarkdownV2 } from "./markdown"; const MAX_ATTACHMENT_BYTES = 50 * 1024 * 1024; const PHOTO_EXTS = new Set([".jpg", ".jpeg", ".png", ".gif", ".webp"]); @@ -219,29 +219,38 @@ export class Outbound { // ---- model-tool helpers ------------------------------------------------ - /** Send text to a chat, chunked + MarkdownV2 (plain fallback on parse error). Returns message ids. */ + /** Send text with the configured Markdown format and safe labeled fallback. Returns actual message ids. */ async send(chatId: string, text: string, opts?: { replyTo?: number; format?: "text" | "markdown"; threadId?: number }): Promise { + if (!text) return []; const access = this.#getAccess(); - const budget = messageLimit(access) - MARKDOWN_HEADROOM; - const parts = chunkLabeled(text, budget, access.chunkMode ?? "newline"); - if (parts.length === 0) return []; const replyMode = access.replyToMode ?? "first"; const useMd = (opts?.format ?? "markdown") === "markdown"; const ids: number[] = []; let threadId = opts?.threadId; let recovered = false; - for (let i = 0; i < parts.length; i++) { - const replyTo = this.#threadTarget(opts?.replyTo, replyMode, i); + const deliver = async (op: () => Promise): Promise => { try { - ids.push(await this.#sendOne(chatId, parts[i], useMd, replyTo, threadId)); + return await op(); } catch (err) { if (recovered || threadId == null || !isMissingThreadError(err)) throw err; const replacement = await this.#recoverMissingThread(chatId, threadId); if (replacement == null) throw err; recovered = true; threadId = replacement; - ids.push(await this.#sendOne(chatId, parts[i], useMd, replyTo, threadId)); + return op(); } + }; + let allowRich = true; + const richBudget = access.textChunkLimit == null ? TELEGRAM_RICH_MAX_CHARS : messageLimit(access); + if (text.length <= richBudget && this.#wantsRich(text, useMd)) { + const id = await deliver(() => this.#tryRichSend(chatId, text, this.#threadTarget(opts?.replyTo, replyMode, 0), threadId)); + if (id !== undefined) return [id]; + allowRich = false; // Re-split the rejected source; never retry rich for these parts. + } + const parts = chunkLabeled(text, messageLimit(access) - MARKDOWN_HEADROOM, access.chunkMode ?? "newline"); + for (let i = 0; i < parts.length; i++) { + const replyTo = this.#threadTarget(opts?.replyTo, replyMode, i); + ids.push(await deliver(() => this.#sendOne(chatId, parts[i], useMd, replyTo, threadId, allowRich))); } return ids; } @@ -391,14 +400,25 @@ export class Outbound { return para > limit / 2 ? para : line > limit / 2 ? line : space > 0 ? space : limit; } - /** Finalize a live preview: MarkdownV2 attempt then plain fallback, cursor removed. */ + /** Finalize a live preview with the selected format, removing the cursor. */ async #finalizePreview(st: ChatState, text: string, useMd: boolean): Promise { if (st.previewMsgId == null) return; await this.#editDelivered(st.chatId, st.previewMsgId, text, useMd); } - /** Edit an already-delivered message: MarkdownV2 attempt then plain fallback. */ + /** Edit in place; only definitive rich rejections may fall back to legacy formatting. */ async #editDelivered(chatId: string, messageId: number, text: string, useMd: boolean): Promise { + if (this.#wantsRich(text, useMd)) { + try { + await this.#rateLimited(() => + tg(this.#token, "editMessageText", { chat_id: chatId, message_id: messageId, rich_message: { markdown: text } }), + ); + return; + } catch (err) { + if (err instanceof TgError && err.code === 400 && /message is not modified/i.test(err.message)) return; + if (isMissingThreadError(err) || !(err instanceof TgError && (err.code === 400 || err.code === 404))) throw err; + } + } if (useMd) { try { await this.#rateLimited(() => @@ -413,34 +433,41 @@ export class Outbound { } /** Finalize one turn into real message(s), then reset per-turn state. */ - async #finalize(st: ChatState, fullText: string, allowRecovery = true): Promise { + async #finalize(st: ChatState, fullText: string, allowRecovery = true, allowRich = true): Promise { if (st.inflight) await st.inflight.catch(() => {}); // barrier: let any in-flight push settle const access = this.#getAccess(); const budget = messageLimit(access) - MARKDOWN_HEADROOM; const mode = access.chunkMode ?? "newline"; const prior = st.committed.length; try { + let richSent = false; + const richBudget = access.textChunkLimit == null ? TELEGRAM_RICH_MAX_CHARS : messageLimit(access); + if (allowRich && st.previewMsgId == null && st.sentUpTo === 0 && prior === 0 && + fullText.length <= richBudget && this.#wantsRich(fullText, true)) { + richSent = (await this.#tryRichSend(st.chatId, fullText, undefined, st.threadId)) !== undefined; + if (!richSent) allowRich = false; + } if (st.previewMsgId != null) { const rest = fullText.slice(st.sentUpTo); const parts = chunkLabeled(rest, budget, mode, prior); await this.#finalizePreview(st, parts[0] ?? rest, true); - for (let i = 1; i < parts.length; i++) await this.#sendOne(st.chatId, parts[i], true, undefined, st.threadId); + for (let i = 1; i < parts.length; i++) await this.#sendOne(st.chatId, parts[i], true, undefined, st.threadId, allowRich); await this.#labelCommitted(st, prior + Math.max(parts.length, 1)); - } else { + } else if (!richSent) { // `sentUpTo` is non-zero when stream overflow already committed a head // message this turn — resending from 0 would duplicate it. const parts = chunkLabeled(fullText.slice(st.sentUpTo), budget, mode, prior); - for (const part of parts) await this.#sendOne(st.chatId, part, true, undefined, st.threadId); + for (const part of parts) await this.#sendOne(st.chatId, part, true, undefined, st.threadId, allowRich); await this.#labelCommitted(st, prior + parts.length); - if (st.draftId != null) { - // Clear the ephemeral draft so it doesn't linger beside the real message. - await tg(this.#token, "sendMessageDraft", { - chat_id: st.chatId, - draft_id: st.draftId, - text: "", - ...(st.threadId != null ? { message_thread_id: st.threadId } : {}), - }).catch(() => {}); - } + } + if (st.draftId != null) { + // Clear the ephemeral draft so it doesn't linger beside the real message. + await tg(this.#token, "sendMessageDraft", { + chat_id: st.chatId, + draft_id: st.draftId, + text: "", + ...(st.threadId != null ? { message_thread_id: st.threadId } : {}), + }).catch(() => {}); } } catch (err) { if (allowRecovery && st.threadId != null && isMissingThreadError(err)) { @@ -451,7 +478,7 @@ export class Outbound { st.draftId = undefined; st.sentUpTo = 0; // the old topic is gone — redeliver the whole answer st.committed = []; - await this.#finalize(st, fullText, false); + await this.#finalize(st, fullText, false, allowRich); return; } } catch (recoveryError) { @@ -495,7 +522,34 @@ export class Outbound { return replacement; } - async #sendOne(chatId: string, text: string, useMd: boolean, replyTo: number | undefined, threadId?: number): Promise { + #wantsRich(text: string, useMd: boolean): boolean { + if (!useMd || !text) return false; + const mode = this.#getAccess().richMessages; + return mode === "on" || (mode === "auto" && hasRichConstructs(text)); + } + + async #tryRichSend(chatId: string, text: string, replyTo: number | undefined, threadId?: number): Promise { + try { + const sent = await this.#rateLimited(() => + tg<{ message_id: number }>(this.#token, "sendRichMessage", { + chat_id: chatId, + rich_message: { markdown: text }, + ...(threadId != null ? { message_thread_id: threadId } : {}), + ...(replyTo != null ? { reply_parameters: { message_id: replyTo } } : {}), + }), + ); + return sent.message_id; + } catch (err) { + if (isMissingThreadError(err) || !(err instanceof TgError && (err.code === 400 || err.code === 404))) throw err; + return undefined; + } + } + + async #sendOne(chatId: string, text: string, useMd: boolean, replyTo: number | undefined, threadId?: number, allowRich = true): Promise { + if (allowRich && this.#wantsRich(text, useMd)) { + const id = await this.#tryRichSend(chatId, text, replyTo, threadId); + if (id !== undefined) return id; + } const reply = replyTo != null ? { reply_parameters: { message_id: replyTo } } : {}; const thread = threadId != null ? { message_thread_id: threadId } : {}; if (useMd) {