feat(delivery): preserve text format intent

This commit is contained in:
xmanrui 2026-08-24 02:33:19 +08:00
parent 73f262711b
commit dd66f5cbe1
5 changed files with 177 additions and 12 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,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) {

View file

@ -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);
});

View file

@ -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/);
});