From dd66f5cbe1a682c8ed22f40b8f062fbd43cfcd28 Mon Sep 17 00:00:00 2001 From: xmanrui <841206367@qq.com> Date: Mon, 24 Aug 2026 02:33:19 +0800 Subject: [PATCH] feat(delivery): preserve text format intent --- .../shared/semantic/artifact-delivery.mjs | 2 +- src/channels/shared/semantic/delivery.mjs | 38 +++++++++ src/channels/shared/text-harness-bridge.mjs | 78 ++++++++++++++++--- test/artifact-delivery.test.mjs | 14 ++++ test/delivery-receipt.test.mjs | 57 ++++++++++++++ 5 files changed, 177 insertions(+), 12 deletions(-) diff --git a/src/channels/shared/semantic/artifact-delivery.mjs b/src/channels/shared/semantic/artifact-delivery.mjs index 7bef533..49868df 100644 --- a/src/channels/shared/semantic/artifact-delivery.mjs +++ b/src/channels/shared/semantic/artifact-delivery.mjs @@ -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; diff --git a/src/channels/shared/semantic/delivery.mjs b/src/channels/shared/semantic/delivery.mjs index 9292686..df187f5 100644 --- a/src/channels/shared/semantic/delivery.mjs +++ b/src/channels/shared/semantic/delivery.mjs @@ -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()], }); } diff --git a/src/channels/shared/text-harness-bridge.mjs b/src/channels/shared/text-harness-bridge.mjs index 251af32..74e505c 100644 --- a/src/channels/shared/text-harness-bridge.mjs +++ b/src/channels/shared/text-harness-bridge.mjs @@ -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,11 +541,18 @@ 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; @@ -529,9 +562,11 @@ export class TextHarnessBridge { // 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 (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 +579,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 +613,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 +626,8 @@ export class TextHarnessBridge { return; } this.#logger.error?.(`[dsh-im:${this.#descriptor.key}] failed to process a message:`, error); + if (await presentStreamFailure('消息处理失败,请稍后重试。')) return; + stream?.cancel?.(); try { await this.#bot.sendText(target, '消息处理失败,请稍后重试。'); } catch (sendError) { diff --git a/test/artifact-delivery.test.mjs b/test/artifact-delivery.test.mjs index 06600d2..7a0d615 100644 --- a/test/artifact-delivery.test.mjs +++ b/test/artifact-delivery.test.mjs @@ -330,3 +330,17 @@ test('mixed artifacts preserve order and merge existing receipt semantics', asyn await assertReleased(image); await assertReleased(file); }); + +test('a definitively failed base delivery is not treated as user-visible', async () => { + const delivery = await deliverOutboundArtifacts({ + baseReceipt: createDeliveryReceipt({ + deliveryId: 'failed-text', + presentation: 'telegram-rich-final', + deliveryOutcome: 'failed', + reason: 'telegram-provider-rejected', + }), + channelKey: 'telegram', + }); + + assert.equal(delivery.userVisible, false); +}); diff --git a/test/delivery-receipt.test.mjs b/test/delivery-receipt.test.mjs index 620a71d..c73b246 100644 --- a/test/delivery-receipt.test.mjs +++ b/test/delivery-receipt.test.mjs @@ -5,10 +5,38 @@ import { artifactOutcomeForError, createArtifactFailureReceipt, createDeliveryReceipt, + createTextDeliveryBlock, mergeDeliveryReceipts, providerMessageIdsFor, } from '../src/channels/shared/semantic/delivery.mjs'; +test('text DeliveryBlocks preserve explicit format and legacy strings default to plain', () => { + const legacy = createTextDeliveryBlock('legacy *text*'); + const markdown = createTextDeliveryBlock({ + kind: 'text', + text: '# Harness answer', + format: 'markdown', + }); + + assert.deepEqual(legacy, { + kind: 'text', + text: 'legacy *text*', + format: 'plain', + }); + assert.deepEqual(markdown, { + kind: 'text', + text: '# Harness answer', + format: 'markdown', + }); + assert.equal(Object.isFrozen(legacy), true); + assert.equal(Object.isFrozen(markdown), true); + assert.throws( + () => createTextDeliveryBlock({ kind: 'text', text: 'x', format: 'html' }), + /plain or markdown/, + ); + assert.throws(() => createTextDeliveryBlock(' '), /non-empty text/); +}); + test('DeliveryReceipt validates and freezes the shared versioned contract', () => { const receipt = createDeliveryReceipt({ deliveryId: 'delivery-1', @@ -102,3 +130,32 @@ test('text and multiple artifact attempts merge into one authoritative receipt', ], }); }); + +test('DeliveryReceipt optionally preserves final delivery outcome without changing legacy receipts', () => { + const legacy = createDeliveryReceipt({ + deliveryId: 'legacy', + presentation: 'telegram-text', + }); + const uncertain = createDeliveryReceipt({ + deliveryId: 'rich', + presentation: 'telegram-rich-final', + deliveryOutcome: 'unknown', + reason: 'telegram-timeout', + }); + + assert.equal(Object.hasOwn(legacy, 'deliveryOutcome'), false); + assert.deepEqual(uncertain, { + schemaVersion: 1, + deliveryId: 'rich', + presentation: 'telegram-rich-final', + providerMessageIds: [], + deliveryOutcome: 'unknown', + reason: 'telegram-timeout', + artifacts: [], + }); + assert.throws(() => createDeliveryReceipt({ + deliveryId: 'invalid', + presentation: 'telegram-rich-final', + deliveryOutcome: 'maybe', + }), /deliveryOutcome/); +});