merge: add native Telegram rich delivery

# Conflicts:
#	lib/index.js
This commit is contained in:
xmanrui 2026-08-24 03:44:19 +08:00
commit a5c516a61c
12 changed files with 1978 additions and 184 deletions

View file

@ -76,7 +76,7 @@ export async function deliverOutboundArtifacts({
logger,
}) {
const receipts = baseReceipt ? [baseReceipt] : [];
let userVisible = Boolean(baseReceipt);
let userVisible = Boolean(baseReceipt) && baseReceipt.deliveryOutcome !== 'failed';
let failureNoticeVisible = false;
let artifactsSent = 0;
let artifactSendErrors = 0;

View file

@ -1,6 +1,8 @@
export const DELIVERY_RECEIPT_SCHEMA_VERSION = 1;
const ARTIFACT_OUTCOMES = new Set(['sent', 'rejected', 'failed', 'unknown']);
const DELIVERY_OUTCOMES = new Set(['sent', 'failed', 'unknown']);
const TEXT_FORMATS = new Set(['plain', 'markdown']);
const REJECTED_ARTIFACT_ERRORS = new Set([
'artifact-changed',
'artifact-context-required',
@ -52,6 +54,24 @@ function requiredString(value, name) {
return value;
}
export function createTextDeliveryBlock(value, format = 'plain') {
const block = typeof value === 'string'
? { kind: 'text', text: value, format }
: value;
if (!block || typeof block !== 'object' || Array.isArray(block)
|| block.kind !== 'text' || typeof block.text !== 'string' || !block.text.trim()) {
throw new TypeError('text delivery block must contain non-empty text');
}
if (!TEXT_FORMATS.has(block.format)) {
throw new TypeError('text delivery format must be plain or markdown');
}
return Object.freeze({
kind: 'text',
text: block.text,
format: block.format,
});
}
function providerIds(values) {
if (!Array.isArray(values)) throw new TypeError('providerMessageIds must be an array');
const ids = [];
@ -98,12 +118,22 @@ export function createDeliveryReceipt({
presentation,
providerMessageIds = [],
artifacts = [],
deliveryOutcome,
reason,
}) {
if (deliveryOutcome !== undefined && !DELIVERY_OUTCOMES.has(deliveryOutcome)) {
throw new TypeError('deliveryOutcome must be sent, failed, or unknown');
}
const normalizedReason = reason === undefined
? undefined
: requiredString(reason, 'delivery reason');
return Object.freeze({
schemaVersion: DELIVERY_RECEIPT_SCHEMA_VERSION,
deliveryId: requiredString(deliveryId, 'deliveryId'),
presentation: requiredString(presentation, 'presentation'),
providerMessageIds: providerIds(providerMessageIds),
...(deliveryOutcome === undefined ? {} : { deliveryOutcome }),
...(normalizedReason === undefined ? {} : { reason: normalizedReason }),
artifacts: artifactResults(artifacts),
});
}
@ -137,11 +167,17 @@ export function mergeDeliveryReceipts({ deliveryId, presentation, receipts }) {
}
const messageIds = [];
const artifacts = new Map();
let deliveryOutcome;
let reason;
for (const receipt of receipts) {
if (!receipt || receipt.schemaVersion !== DELIVERY_RECEIPT_SCHEMA_VERSION) {
throw new TypeError('receipt must use DeliveryReceipt schema version 1');
}
messageIds.push(...(receipt.providerMessageIds ?? []));
if (deliveryOutcome === undefined && receipt.deliveryOutcome !== undefined) {
deliveryOutcome = receipt.deliveryOutcome;
reason = receipt.reason;
}
for (const artifact of receipt.artifacts ?? []) {
artifacts.set(artifact.artifactId, artifact);
}
@ -150,6 +186,8 @@ export function mergeDeliveryReceipts({ deliveryId, presentation, receipts }) {
deliveryId,
presentation,
providerMessageIds: messageIds,
deliveryOutcome,
reason,
artifacts: [...artifacts.values()],
});
}

View file

@ -35,6 +35,7 @@ import {
import { deliverOutboundArtifacts } from './semantic/artifact-delivery.mjs';
import {
createDeliveryReceipt,
createTextDeliveryBlock,
providerMessageIdsFor,
} from './semantic/delivery.mjs';
@ -372,6 +373,7 @@ export class TextHarnessBridge {
const target = message.replyTarget;
const text = cleanText(message.content);
let stream = null;
let semanticStream = false;
try {
this.#signal?.throwIfAborted();
if (message.kind === 'group' && message.addressed !== true) {
@ -448,7 +450,17 @@ export class TextHarnessBridge {
this.#logger.warn?.(`[dsh-im:${this.#descriptor.key}] typing indicator failed:`, error);
});
let streamFinished = false;
if (typeof this.#bot.openStream === 'function') {
if (typeof this.#bot.openDeliveryStream === 'function') {
try {
stream = await this.#bot.openDeliveryStream(target);
semanticStream = true;
} catch (error) {
this.#logger.warn?.(
`[dsh-im:${this.#descriptor.key}] unable to start a semantic reply stream; using final delivery:`,
error,
);
}
} else if (typeof this.#bot.openStream === 'function') {
try {
stream = await this.#bot.openStream(target);
} catch (error) {
@ -476,7 +488,12 @@ export class TextHarnessBridge {
onUpdate: stream ? async (update) => {
const progress = update.type === 'text' ? update.text
: update.type === 'tool' ? `正在使用${update.name}…` : update.text;
if (progress) await stream.update(progress);
if (progress) {
const format = update.type === 'text' ? 'markdown' : 'plain';
await stream.update(semanticStream
? createTextDeliveryBlock(progress, format)
: progress);
}
} : undefined,
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key: conversationKey,
@ -491,19 +508,28 @@ export class TextHarnessBridge {
const visibleAnswer = !cleanText(answer) && artifacts.length > 0
? FILE_ONLY_COMPLETION_TEXT
: answer;
const answerFormat = visibleAnswer === FILE_ONLY_COMPLETION_TEXT && !cleanText(answer)
? 'plain'
: 'markdown';
let textDeliveryError = null;
let textReceipt = null;
if (stream) {
try {
const result = await stream.finish(visibleAnswer);
const result = await stream.finish(semanticStream
? createTextDeliveryBlock(visibleAnswer, answerFormat)
: visibleAnswer);
streamFinished = true;
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: `${this.#descriptor.key}-stream`,
presentation: result?.presentation
?? stream.presentation
?? `${this.#descriptor.key}-stream`,
providerMessageIds: [
...providerMessageIdsFor(stream),
...providerMessageIdsFor(result),
],
deliveryOutcome: result?.deliveryOutcome,
reason: result?.reason,
});
} catch (error) {
stream.cancel?.();
@ -515,23 +541,40 @@ export class TextHarnessBridge {
}
if (!streamFinished) {
try {
const result = await this.#bot.sendText(target, visibleAnswer);
const result = typeof this.#bot.sendDelivery === 'function'
? await this.#bot.sendDelivery(
target,
createTextDeliveryBlock(visibleAnswer, answerFormat),
)
: await this.#bot.sendText(target, visibleAnswer);
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: `${this.#descriptor.key}-text`,
presentation: result?.presentation ?? `${this.#descriptor.key}-text`,
providerMessageIds: providerMessageIdsFor(result),
deliveryOutcome: result?.deliveryOutcome,
reason: result?.reason,
});
} catch (error) {
textDeliveryError = error;
}
}
if (textReceipt?.deliveryOutcome === 'failed') {
const reason = textReceipt.reason ?? 'text-delivery-failed';
textDeliveryError = new Error(`Final text delivery failed (${reason})`);
textDeliveryError.code = reason;
}
// A failed final text must not discard an already registered result file.
// Settle the independent attachment path before surfacing the text error.
const delivery = await this.#deliverArtifacts(target, messageId, artifacts, textReceipt);
if (textDeliveryError && !delivery.userVisible) throw textDeliveryError;
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
if (textDeliveryError && !delivery.userVisible) {
textDeliveryError.deliveryReceipt = delivery.receipt;
throw textDeliveryError;
}
if (delivery.userVisible) {
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
}
return delivery.receipt;
} catch (error) {
if (error?.code === 'turn-stopped') {
@ -544,11 +587,28 @@ export class TextHarnessBridge {
}
return;
}
stream?.cancel?.();
if (this.#signal?.aborted) return;
if (this.#signal?.aborted) {
stream?.cancel?.();
return;
}
this.#status.lastError = error?.message ?? String(error);
const presentStreamFailure = async (text) => {
if (typeof stream?.fail !== 'function') return false;
try {
const result = await stream.fail(text);
return Boolean(result) && result.deliveryOutcome !== 'failed';
} catch (streamError) {
this.#logger.warn?.(
`[dsh-im:${this.#descriptor.key}] unable to finalize the failed stream:`,
streamError,
);
return false;
}
};
const imageErrorMessage = imagePromptUserMessage(error);
if (imageErrorMessage) {
if (await presentStreamFailure(imageErrorMessage)) return;
stream?.cancel?.();
try {
await this.#bot.sendText(target, imageErrorMessage);
} catch (sendError) {
@ -561,6 +621,8 @@ export class TextHarnessBridge {
}
const fileErrorMessage = inboundFileUserMessage(error);
if (fileErrorMessage) {
if (await presentStreamFailure(fileErrorMessage)) return;
stream?.cancel?.();
try {
await this.#bot.sendText(target, fileErrorMessage);
} catch (sendError) {
@ -572,6 +634,10 @@ export class TextHarnessBridge {
return;
}
this.#logger.error?.(`[dsh-im:${this.#descriptor.key}] failed to process a message:`, error);
if (await presentStreamFailure('消息处理失败,请稍后重试。')) {
return error.deliveryReceipt;
}
stream?.cancel?.();
try {
await this.#bot.sendText(target, '消息处理失败,请稍后重试。');
} catch (sendError) {
@ -580,6 +646,7 @@ export class TextHarnessBridge {
sendError,
);
}
return error.deliveryReceipt;
} finally {
await Promise.allSettled([
this.#cancelPendingInteraction(conversationKey),

View file

@ -35,6 +35,25 @@ function preserveProviderMetadata(target, source) {
return target;
}
function inputRichMessage(value) {
if (!value || typeof value !== 'object' || Array.isArray(value)) {
throw new TypeError('A Telegram rich message is required');
}
const formats = ['markdown', 'html', 'blocks'].filter((key) => value[key] !== undefined);
if (formats.length !== 1) {
throw new TypeError('A Telegram rich message requires exactly one format');
}
const selected = formats[0];
if ((selected === 'markdown' || selected === 'html')
&& (typeof value[selected] !== 'string' || !value[selected].trim())) {
throw new TypeError('Telegram rich message text is invalid');
}
if (selected === 'blocks' && !Array.isArray(value.blocks)) {
throw new TypeError('Telegram rich message blocks are invalid');
}
return value;
}
function telegramArtifactProviderError(cause, mediaLabel = 'document') {
const providerCode = Number(cause?.providerCode);
const status = Number(cause?.status);
@ -183,6 +202,35 @@ export class TelegramApi {
}, { signal });
}
async sendRichMessage({
chatId,
richMessage,
replyToMessageId,
messageThreadId,
signal,
}) {
return this.#call('sendRichMessage', {
chat_id: chatId,
rich_message: inputRichMessage(richMessage),
...(replyToMessageId ? {
reply_parameters: { message_id: replyToMessageId, allow_sending_without_reply: true },
} : {}),
...(messageThreadId ? { message_thread_id: messageThreadId } : {}),
}, { signal });
}
async sendRichMessageDraft({ chatId, draftId, richMessage, messageThreadId, signal }) {
if (!Number.isSafeInteger(draftId) || draftId === 0) {
throw new TypeError('Telegram rich message draft id must be a non-zero integer');
}
return this.#call('sendRichMessageDraft', {
chat_id: chatId,
draft_id: draftId,
rich_message: inputRichMessage(richMessage),
...(messageThreadId ? { message_thread_id: messageThreadId } : {}),
}, { signal });
}
async sendDocument({ chatId, file, replyToMessageId, messageThreadId, signal }) {
return this.#sendArtifact('sendDocument', 'document', 'document', {
chatId,
@ -239,6 +287,9 @@ export class TelegramApi {
});
} catch (error) {
if (signal?.aborted) throw abortReason(signal);
if (error?.deliveryOutcome === 'unknown') {
throw uncertainTelegramDelivery(error?.cause ?? error, mediaLabel);
}
if (error?.code?.startsWith?.('telegram-')) {
throw telegramArtifactProviderError(error, mediaLabel);
}
@ -246,12 +297,17 @@ export class TelegramApi {
}
}
async editMessageText({ chatId, messageId, text, signal }) {
async editMessageText({ chatId, messageId, text, richMessage, signal }) {
if ((text === undefined) === (richMessage === undefined)) {
throw new TypeError('Telegram message edit requires exactly one of text or richMessage');
}
return this.#call('editMessageText', {
chat_id: chatId,
message_id: messageId,
text,
link_preview_options: { is_disabled: true },
...(text === undefined ? { rich_message: inputRichMessage(richMessage) } : {
text,
link_preview_options: { is_disabled: true },
}),
}, { signal });
}
@ -283,6 +339,7 @@ export class TelegramApi {
}
async #call(method, payload, { signal, timeoutMs = 15_000, multipart = false } = {}) {
if (signal?.aborted) throw abortReason(signal);
const url = new URL(this.#baseUrl);
url.pathname = `${url.pathname.replace(/\/$/, '')}/bot${this.#token}/${method}`;
let response;
@ -295,8 +352,14 @@ export class TelegramApi {
redirect: 'error',
});
} catch (error) {
if (error?.name === 'AbortError' || error?.name === 'TimeoutError') throw error;
throw new Error(`Telegram ${method} transport failed`);
const transport = new Error(`Telegram ${method} transport failed`, { cause: error });
transport.code = error?.name === 'TimeoutError'
? 'telegram-timeout'
: error?.name === 'AbortError'
? 'telegram-aborted-after-dispatch'
: 'telegram-transport-error';
transport.deliveryOutcome = 'unknown';
throw transport;
}
let body;
try {
@ -304,6 +367,8 @@ export class TelegramApi {
} catch {
const error = new Error(`Telegram ${method} returned invalid JSON`);
error.status = response?.status;
error.code = 'telegram-response-invalid';
error.deliveryOutcome = 'unknown';
throw error;
}
if (!response.ok || body?.ok !== true) {

View file

@ -0,0 +1,147 @@
export const TELEGRAM_RICH_TEXT_LIMIT = 30_000;
export const TELEGRAM_REGULAR_TEXT_LIMIT = 4_000;
function textValue(value) {
if (typeof value !== 'string' || !value.trim()) {
throw new TypeError('Telegram message text must be a non-empty string');
}
return value;
}
function escapedCharacterLength(value) {
if (value === '&') return 5;
if (value === '<' || value === '>') return 4;
return 1;
}
function escapedLength(value) {
return Array.from(value).reduce((total, character) => (
total + escapedCharacterLength(character)
), 0);
}
function preferredCut(points, offset, hardEnd, limit) {
const minimum = offset + Math.floor(limit * 0.5);
for (let index = hardEnd - 1; index >= minimum; index -= 1) {
if (points[index] === '\n' || points[index] === ' ') return index + 1;
}
return hardEnd;
}
function splitText(value, limit, characterLength = () => 1) {
if (typeof value !== 'string') throw new TypeError('Telegram message text must be a string');
if (!value) return [];
const points = Array.from(value);
const chunks = [];
let offset = 0;
while (offset < points.length) {
let end = offset;
let length = 0;
while (end < points.length) {
const nextLength = characterLength(points[end]);
if (length + nextLength > limit) break;
length += nextLength;
end += 1;
}
if (end === offset) throw new Error('Telegram message character exceeds the chunk limit');
if (end < points.length) end = preferredCut(points, offset, end, limit);
chunks.push(points.slice(offset, end).join(''));
offset = end;
}
return chunks;
}
function fenceMarkers(markdown) {
return [...markdown.matchAll(/^```[^\n]*(?:\n|$)/gm)];
}
function assertCompleteFences(markdown) {
const markers = fenceMarkers(markdown);
if (markers.length % 2 !== 0) {
throw new Error('Telegram Rich Markdown contains an unfinished code fence');
}
return markers;
}
function escapedMarkdown(value) {
return value
.replaceAll('&', '&amp;')
.replaceAll('<', '&lt;')
.replaceAll('>', '&gt;');
}
function plainRichChunks(source, limit) {
return splitText(source, limit, escapedCharacterLength).map((chunk) => ({
source: chunk,
markdown: escapedMarkdown(chunk),
}));
}
function fencedRichChunks(source, opening, closing, limit) {
const body = source.slice(opening.length, source.length - closing.length);
const wrapperLength = escapedLength(opening)
+ Math.max(escapedLength(closing), escapedLength('```'))
+ 1;
const bodyLimit = limit - wrapperLength;
if (bodyLimit < 1) throw new Error('Telegram Rich Markdown code fence exceeds the limit');
const bodyChunks = splitText(body, bodyLimit, escapedCharacterLength);
return bodyChunks.map((bodyChunk, index) => {
const last = index === bodyChunks.length - 1;
const separator = bodyChunk.endsWith('\n') ? '' : '\n';
const renderedClosing = last ? closing : '```';
const sourceChunk = `${index === 0 ? opening : ''}${bodyChunk}${last ? closing : ''}`;
const markdown = escapedMarkdown(`${opening}${bodyChunk}${separator}${renderedClosing}`);
if (Array.from(markdown).length > limit) {
throw new Error('Telegram Rich Markdown code block exceeds the limit');
}
return { source: sourceChunk, markdown };
});
}
export function toTelegramRichMarkdown(value) {
const markdown = textValue(value);
assertCompleteFences(markdown);
// Rich Markdown also accepts HTML tags. Escape entities first so encoded or
// raw model output cannot be interpreted as Telegram-specific markup.
return escapedMarkdown(markdown);
}
export function splitTelegramRichMarkdown(value, limit = TELEGRAM_RICH_TEXT_LIMIT) {
if (!Number.isInteger(limit) || limit < 128 || limit > 32_768) {
throw new TypeError('Telegram Rich Markdown limit is invalid');
}
const source = textValue(value);
const markers = assertCompleteFences(source);
const complete = escapedMarkdown(source);
if (Array.from(complete).length <= limit) return [{ source, markdown: complete }];
const chunks = [];
let cursor = 0;
for (let index = 0; index < markers.length; index += 2) {
const opening = markers[index];
const closing = markers[index + 1];
if (opening.index > cursor) {
chunks.push(...plainRichChunks(source.slice(cursor, opening.index), limit));
}
const end = closing.index + closing[0].length;
const fenced = source.slice(opening.index, end);
const rendered = escapedMarkdown(fenced);
if (Array.from(rendered).length <= limit) {
chunks.push({ source: fenced, markdown: rendered });
} else {
chunks.push(...fencedRichChunks(fenced, opening[0], closing[0], limit));
}
cursor = end;
}
if (cursor < source.length) chunks.push(...plainRichChunks(source.slice(cursor), limit));
return chunks;
}
export function splitTelegramRegularText(value, limit = TELEGRAM_REGULAR_TEXT_LIMIT) {
if (!Number.isInteger(limit) || limit < 1 || limit > 4_096) {
throw new TypeError('Telegram regular message limit is invalid');
}
// Keep Telegram's established 4000 UTF-16-unit boundary without ever cutting
// a surrogate pair in half.
return splitText(textValue(value), limit, (character) => character.length);
}

View file

@ -1,6 +1,14 @@
import { randomInt } from 'node:crypto';
import { createEditableMessageStream, splitMessageText } from '../shared/editable-message-stream.mjs';
import { createTextDeliveryBlock } from '../shared/semantic/delivery.mjs';
import { COMMANDS_MENU_BUTTON, TelegramApi } from './telegram-api.mjs';
import { createTelegramBridgeStatus, TelegramHarnessBridge } from './telegram-bridge.mjs';
import {
splitTelegramRegularText,
splitTelegramRichMarkdown,
toTelegramRichMarkdown,
} from './telegram-rich-message.mjs';
import {
TELEGRAM_ACCESS_MODES,
normalizeTelegramAccessPolicy,
@ -149,6 +157,7 @@ export function normalizeTelegramUpdate(update, {
addressed,
replyTarget: {
chatId,
chatType: message.chat.type,
replyToMessageId: messageId,
messageThreadId,
},
@ -166,17 +175,105 @@ export function telegramInboundAllowed(message, {
&& allowedPrivateUserIds.has(String(message.senderId));
}
function telegramMessageId(value) {
return Number.isSafeInteger(value?.message_id) ? String(value.message_id) : null;
}
function telegramFailure(error) {
const providerCode = Number(error?.providerCode);
const status = Number(error?.status);
const unknown = error?.deliveryOutcome === 'unknown'
|| providerCode >= 500
|| status >= 500
|| ['telegram-timeout', 'telegram-aborted-after-dispatch', 'telegram-transport-error',
'telegram-response-invalid'].includes(error?.code);
return {
outcome: unknown ? 'unknown' : 'failed',
reason: typeof error?.code === 'string' && error.code
? error.code
: unknown ? 'telegram-delivery-uncertain' : 'telegram-provider-rejected',
};
}
function deliveryResult(presentation, providerMessageIds, deliveryOutcome = 'sent', reason) {
return {
presentation,
providerMessageIds: [...new Set(providerMessageIds)],
deliveryOutcome,
...(reason ? { reason } : {}),
};
}
class TelegramDeliveryStream {
#update;
#finish;
#fail;
#logger;
#providerMessageIds;
#closed = false;
#lastUpdate = null;
constructor({ update, finish, fail, providerMessageIds = [], presentation, logger }) {
this.#update = update;
this.#finish = finish;
this.#fail = fail;
this.#providerMessageIds = providerMessageIds;
this.presentation = presentation;
this.#logger = logger;
}
get providerMessageIds() {
return [...this.#providerMessageIds];
}
async update(value) {
if (this.#closed) return undefined;
const block = createTextDeliveryBlock(value);
const key = `${block.format}:${block.text}`;
if (key === this.#lastUpdate) return undefined;
this.#lastUpdate = key;
try {
return await this.#update(block);
} catch (error) {
this.#logger.warn?.('[dsh-im:telegram] rich stream update failed:', error);
return undefined;
}
}
async finish(value) {
if (this.#closed) throw new Error('Message stream is already closed');
this.#closed = true;
const result = await this.#finish(createTextDeliveryBlock(value));
this.#providerMessageIds.push(...(result?.providerMessageIds ?? []));
return result;
}
async fail(text) {
if (this.#closed) return undefined;
this.#closed = true;
const result = await this.#fail(createTextDeliveryBlock(text, 'plain'));
this.#providerMessageIds.push(...(result?.providerMessageIds ?? []));
return result;
}
cancel() {
this.#closed = true;
}
}
export class TelegramBotClient {
#api;
#signal;
#logger;
constructor({ api, signal }) {
constructor({ api, signal, logger = console }) {
this.#api = api;
this.#signal = signal;
this.#logger = logger;
}
async sendText(target, text) {
const chunks = splitMessageText(text, 4_000);
const chunks = splitTelegramRegularText(text);
const providerMessageIds = [];
for (const [index, chunk] of chunks.entries()) {
const result = await this.#api.sendMessage({
@ -221,6 +318,265 @@ export class TelegramBotClient {
});
}
async #sendPlain(target, text, { placeholderMessageId } = {}) {
const chunks = splitTelegramRegularText(text);
const providerMessageIds = placeholderMessageId === undefined
? []
: [String(placeholderMessageId)];
let firstUnsent = 0;
if (placeholderMessageId !== undefined) {
try {
await this.#api.editMessageText({
chatId: target.chatId,
messageId: placeholderMessageId,
text: chunks[0],
signal: this.#signal,
});
firstUnsent = 1;
} catch (error) {
const failure = telegramFailure(error);
if (failure.outcome === 'unknown') {
return deliveryResult('text-fallback', providerMessageIds, 'unknown', failure.reason);
}
const fallback = await this.#sendPlain(target, text);
const terminalText = fallback.deliveryOutcome === 'sent'
? '回复已发送。'
: fallback.deliveryOutcome === 'unknown'
? '回复发送结果未能确认。'
: '消息发送失败,请稍后重试。';
try {
await this.#api.editMessageText({
chatId: target.chatId,
messageId: placeholderMessageId,
text: terminalText,
signal: this.#signal,
});
} catch (terminalError) {
this.#logger.warn?.(
'[dsh-im:telegram] unable to replace the processing placeholder:',
terminalError,
);
}
return deliveryResult(
'text-fallback',
[...providerMessageIds, ...fallback.providerMessageIds],
fallback.deliveryOutcome,
fallback.reason,
);
}
}
for (let index = firstUnsent; index < chunks.length; index += 1) {
try {
const result = await this.#api.sendMessage({
chatId: target.chatId,
text: chunks[index],
replyToMessageId: index === 0 ? target.replyToMessageId : undefined,
messageThreadId: target.messageThreadId,
signal: this.#signal,
});
const id = telegramMessageId(result);
if (id) providerMessageIds.push(id);
} catch (error) {
const failure = telegramFailure(error);
return deliveryResult(
'text-fallback',
providerMessageIds,
failure.outcome,
failure.reason,
);
}
}
return deliveryResult('text-fallback', providerMessageIds);
}
async #sendRich(target, block) {
if (block.format === 'plain') return this.#sendPlain(target, block.text);
let chunks;
try {
chunks = splitTelegramRichMarkdown(block.text);
} catch {
return this.#sendPlain(target, block.text);
}
const providerMessageIds = [];
for (const [index, chunk] of chunks.entries()) {
try {
const result = await this.#api.sendRichMessage({
chatId: target.chatId,
richMessage: { markdown: chunk.markdown },
replyToMessageId: index === 0 ? target.replyToMessageId : undefined,
messageThreadId: target.messageThreadId,
signal: this.#signal,
});
const id = telegramMessageId(result);
if (id) providerMessageIds.push(id);
} catch (error) {
const failure = telegramFailure(error);
if (failure.outcome === 'unknown') {
return deliveryResult(
'telegram-rich-final',
providerMessageIds,
'unknown',
failure.reason,
);
}
const remaining = chunks.slice(index).map((part) => part.source).join('');
const fallback = await this.#sendPlain({
...target,
replyToMessageId: index === 0 ? target.replyToMessageId : undefined,
}, remaining);
return deliveryResult(
fallback.presentation,
[...providerMessageIds, ...fallback.providerMessageIds],
fallback.deliveryOutcome,
fallback.reason,
);
}
}
return deliveryResult('telegram-rich-final', providerMessageIds);
}
async #editRich(target, messageId, block) {
if (block.format === 'plain') {
return this.#sendPlain(target, block.text, {
placeholderMessageId: messageId,
});
}
let chunks;
try {
chunks = splitTelegramRichMarkdown(block.text);
} catch {
return this.#sendPlain(target, block.text, { placeholderMessageId: messageId });
}
const providerMessageIds = [String(messageId)];
try {
await this.#api.editMessageText({
chatId: target.chatId,
messageId,
richMessage: { markdown: chunks[0].markdown },
signal: this.#signal,
});
} catch (error) {
const failure = telegramFailure(error);
if (failure.outcome === 'unknown') {
return deliveryResult(
'telegram-rich-final',
providerMessageIds,
'unknown',
failure.reason,
);
}
return this.#sendPlain(target, block.text, { placeholderMessageId: messageId });
}
for (const [offset, chunk] of chunks.slice(1).entries()) {
try {
const result = await this.#api.sendRichMessage({
chatId: target.chatId,
richMessage: { markdown: chunk.markdown },
messageThreadId: target.messageThreadId,
signal: this.#signal,
});
const id = telegramMessageId(result);
if (id) providerMessageIds.push(id);
} catch (error) {
const failure = telegramFailure(error);
if (failure.outcome === 'unknown') {
return deliveryResult(
'telegram-rich-final',
providerMessageIds,
'unknown',
failure.reason,
);
}
const remaining = chunks.slice(offset + 1).map((part) => part.source).join('');
const fallback = await this.#sendPlain({
...target,
replyToMessageId: undefined,
}, remaining);
return deliveryResult(
fallback.presentation,
[...providerMessageIds, ...fallback.providerMessageIds],
fallback.deliveryOutcome,
fallback.reason,
);
}
}
return deliveryResult('telegram-rich-final', providerMessageIds);
}
sendDelivery(target, value) {
return this.#sendRich(target, createTextDeliveryBlock(value));
}
async openDeliveryStream(target) {
if (target.chatType === 'private') {
const draftId = randomInt(1, 2_147_483_647);
const updateDraft = async (block) => {
const richMessage = { markdown: toTelegramRichMarkdown(block.text) };
await this.#api.sendRichMessageDraft({
chatId: target.chatId,
draftId,
richMessage,
messageThreadId: target.messageThreadId,
signal: this.#signal,
});
return deliveryResult('telegram-rich-draft', []);
};
const stream = new TelegramDeliveryStream({
update: updateDraft,
finish: (block) => this.#sendRich(target, block),
fail: (block) => this.#sendPlain(target, block.text),
presentation: 'telegram-rich-draft',
logger: this.#logger,
});
await stream.update(createTextDeliveryBlock('正在处理…', 'plain'));
return stream;
}
const placeholder = await this.#api.sendMessage({
chatId: target.chatId,
text: '正在处理…',
replyToMessageId: target.replyToMessageId,
messageThreadId: target.messageThreadId,
signal: this.#signal,
});
const messageId = placeholder?.message_id;
if (!Number.isSafeInteger(messageId)) {
throw new Error('Telegram did not return a placeholder message id');
}
return new TelegramDeliveryStream({
update: async (block) => {
if (block.format === 'plain') {
await this.#api.editMessageText({
chatId: target.chatId,
messageId,
text: block.text,
signal: this.#signal,
});
return deliveryResult('text-fallback', [String(messageId)]);
}
await this.#api.editMessageText({
chatId: target.chatId,
messageId,
richMessage: { markdown: toTelegramRichMarkdown(block.text) },
signal: this.#signal,
});
return deliveryResult('telegram-rich-draft', [String(messageId)]);
},
finish: (block) => this.#editRich(target, messageId, block),
fail: (block) => this.#sendPlain(target, block.text, {
placeholderMessageId: messageId,
}),
providerMessageIds: [String(messageId)],
presentation: 'telegram-regular',
logger: this.#logger,
});
}
async openStream(target) {
const stream = createEditableMessageStream({
limit: 4_000,
@ -360,7 +716,11 @@ export class TelegramRuntime {
error,
);
}
const client = new TelegramBotClient({ api, signal: controller.signal });
const client = new TelegramBotClient({
api,
signal: controller.signal,
logger: this.#logger,
});
this.#bridge = new TelegramHarnessBridge({
bot: client,
harness: this.#harness,