feat(telegram): deliver native rich messages

This commit is contained in:
xmanrui 2026-08-24 02:49:16 +08:00
parent dd66f5cbe1
commit d7f4cf03ab
7 changed files with 1588 additions and 170 deletions

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,