fix: preserve quoted message context across channels

This commit is contained in:
xmanrui 2026-09-02 01:27:35 +08:00
parent c2be2389b0
commit 94e114b9b3
46 changed files with 4242 additions and 321 deletions

View file

@ -39,7 +39,6 @@ import {
hasInboundImages,
imagePromptDiagnostic,
imagePromptUserMessage,
promptContentForMessage,
} from '../shared/image-prompt.mjs';
import {
hasInboundFiles,
@ -48,10 +47,16 @@ import {
} from '../shared/inbound-file.mjs';
import { rememberConnectionTestTarget } from '../shared/connection-test.mjs';
import { deliverOutboundArtifacts } from '../shared/semantic/artifact-delivery.mjs';
import {
hasReplyReference,
promptContentForInboundMessage,
} from '../shared/semantic/reply-reference.mjs';
import {
createDeliveryReceipt,
providerMessageIdsFor,
} from '../shared/semantic/delivery.mjs';
import { recoverAssistantTextByTimestamp } from '../shared/session-reply-recovery.mjs';
import { DINGTALK_RECENT_OUTBOUND_MATCH_TOLERANCE_MS } from './state-store.mjs';
import {
channelDeliveryFailure,
clearLastMessageFailure,
@ -185,11 +190,95 @@ function downloadCodeFor(value) {
return nonEmptyString(value?.downloadCode) ?? nonEmptyString(value?.pictureDownloadCode);
}
function dingtalkTimestampMs(value) {
const number = typeof value === 'string' && value.trim() ? Number(value) : value;
if (!Number.isFinite(number) || number < 0) return null;
return Math.trunc(number < 10_000_000_000 ? number * 1_000 : number);
}
function usefulReplyText(value) {
const text = nonEmptyString(value);
return text && !/^\[interactive card message\]$/iu.test(text) ? text : null;
}
function dingtalkReplyReference(message, options) {
const replyEnvelope = message?.text;
if (replyEnvelope?.isReplyMsg !== true) return null;
const replied = replyEnvelope?.repliedMsg;
if (!replied || typeof replied !== 'object') {
return { unavailableReason: 'not-delivered' };
}
const msgtype = nonEmptyString(replied.msgType ?? replied.msgtype)?.toLowerCase() ?? '';
const repliedContent = parsedMessageContent({ content: replied.content }) ?? {};
const pseudoMessage = {
msgtype,
text: {
content: nonEmptyString(repliedContent.text)
?? (typeof replied.content === 'string' ? replied.content : ''),
},
content: repliedContent,
};
const normalized = dingtalkInboundMessage(pseudoMessage, options);
let attachments = [];
if (msgtype === 'picture') {
attachments = [{ kind: 'image' }];
} else if (msgtype === 'file') {
const name = nonEmptyString(repliedContent.fileName ?? repliedContent.file_name);
attachments = [{ kind: 'file', ...(name ? { name } : {}) }];
} else if (msgtype === 'richtext') {
attachments = richTextEntries(repliedContent)
.filter((entry) => String(entry?.type ?? '').toLowerCase() === 'picture')
.map(() => ({ kind: 'image' }));
} else if (msgtype === 'voice' || msgtype === 'audio') {
attachments = [{ kind: 'audio' }];
} else if (msgtype === 'video') {
attachments = [{ kind: 'video' }];
}
const messageId = nonEmptyString(replied.msgId ?? replied.messageId);
const authorId = nonEmptyString(replied.senderId ?? replied.senderStaffId);
const authorName = nonEmptyString(replied.senderNick ?? replied.senderName);
const content = usefulReplyText(normalized.content)
?? usefulReplyText(repliedContent.text)
?? usefulReplyText(repliedContent.summary)
?? usefulReplyText(repliedContent.title);
const processQueryKey = nonEmptyString(
message?.originalProcessQueryKey ?? repliedContent.processQueryKey,
);
const createdAt = dingtalkTimestampMs(replied.createdAt ?? replied.createTime);
const load = !content && attachments.length === 0
&& typeof options?.loadReplyContent === 'function'
? ({ signal } = {}) => options.loadReplyContent({
...(messageId ? { messageId } : {}),
...(processQueryKey ? { processQueryKey } : {}),
...(createdAt === null ? {} : { createdAt }),
}, { signal })
: null;
const supported = [
'text', 'picture', 'file', 'richtext', 'voice', 'audio', 'video',
'interactivecard', 'chatrecord',
]
.includes(msgtype);
return {
...(messageId ? { messageId } : {}),
...(authorId ? { authorId } : {}),
...(authorName ? { authorName } : {}),
...(content ? { content } : {}),
...(attachments.length > 0 ? { attachments } : {}),
...(load ? { load } : {}),
...(!content && attachments.length === 0 && !load
? { unavailableReason: supported ? 'not-delivered' : 'unsupported' }
: {}),
};
}
/** Normalize DingTalk picture and richText callbacks into lazy image references. */
export function dingtalkInboundMessage(message, {
api,
clientId,
clientSecret,
loadReplyContent,
} = {}) {
const msgtype = String(message?.msgtype ?? '').toLowerCase();
const content = parsedMessageContent(message);
@ -209,6 +298,12 @@ export function dingtalkInboundMessage(message, {
}
}
const fileCode = msgtype === 'file' ? downloadCodeFor(content) : null;
const replyTo = dingtalkReplyReference(message, {
api,
clientId,
clientSecret,
loadReplyContent,
});
return {
content: text,
images: imageCodes.map((downloadCode, index) => ({
@ -242,6 +337,7 @@ export function dingtalkInboundMessage(message, {
});
},
}] : [],
...(replyTo ? { replyTo } : {}),
};
}
@ -465,11 +561,7 @@ export class DingtalkHarnessBridge {
// An unsafe reply route must never be able to submit an approval.
}
const pending = this.#pendingInteractions.get(key);
const promptMessage = dingtalkInboundMessage(message, {
api: this.#api,
clientId: this.#clientId,
clientSecret: this.#clientSecret,
});
const promptMessage = this.#inboundMessage(message, key);
const commandText = nonEmptyString(promptMessage.content) ?? '';
const addressed = String(message.conversationType) !== '2' || message?.isInAtList === true;
const direct = String(message.conversationType) !== '2';
@ -533,7 +625,8 @@ export class DingtalkHarnessBridge {
plainText: Boolean(commandText)
&& String(message?.msgtype).toLowerCase() === 'text'
&& !hasInboundFiles(promptMessage)
&& !hasInboundImages(promptMessage),
&& !hasInboundImages(promptMessage)
&& !hasReplyReference(promptMessage),
});
if (result.handled) {
if (result.kind === 'submit') {
@ -778,11 +871,7 @@ export class DingtalkHarnessBridge {
}
const addressed = String(message.conversationType) !== '2' || message.isInAtList === true;
const preparedMessage = hasSafeReplyRoute && addressed
? prefetchInboundFiles(dingtalkInboundMessage(message, {
api: this.#api,
clientId: this.#clientId,
clientSecret: this.#clientSecret,
}), { signal: this.#signal })
? prefetchInboundFiles(this.#inboundMessage(message, key), { signal: this.#signal })
: undefined;
const previous = this.#queues.get(key) ?? Promise.resolve();
const current = previous
@ -801,6 +890,50 @@ export class DingtalkHarnessBridge {
return current;
}
#inboundMessage(message, key) {
return dingtalkInboundMessage(message, {
api: this.#api,
clientId: this.#clientId,
clientSecret: this.#clientSecret,
loadReplyContent: (reference, options) => this.#loadReplyContent(key, reference, options),
});
}
async #loadReplyContent(key, reference, { signal } = {}) {
const indexed = this.#state.recentOutboundTextFor?.({
conversationKey: key,
...reference,
});
if (indexed) return { content: indexed };
const quotedAt = dingtalkTimestampMs(reference?.createdAt);
if (quotedAt === null) return { unavailableReason: 'not-delivered' };
const sessionId = this.#state.sessionFor(key);
const session = typeof sessionId === 'string' && sessionId
? this.#harness.workspaceSession?.(sessionId)
: null;
const text = await recoverAssistantTextByTimestamp({
session,
quotedAt,
signal,
toleranceMs: DINGTALK_RECENT_OUTBOUND_MATCH_TOLERANCE_MS,
});
if (!text) return { unavailableReason: 'not-delivered' };
try {
await this.#state.rememberOutboundMessage?.({
conversationKey: key,
text,
sentAt: quotedAt,
completedAt: quotedAt,
providerMessageIds: [reference?.processQueryKey, reference?.messageId]
.map(nonEmptyString)
.filter(Boolean),
});
} catch (error) {
this.#logger.warn?.('[dsh-dingtalk] failed to remember a recovered quote:', error);
}
return { content: text };
}
async waitForIdle() {
await Promise.allSettled([
...this.#queues.values(),
@ -936,20 +1069,18 @@ export class DingtalkHarnessBridge {
return;
}
const promptMessage = preparedMessage ?? dingtalkInboundMessage(message, {
api: this.#api,
clientId: this.#clientId,
clientSecret: this.#clientSecret,
});
const promptMessage = preparedMessage ?? this.#inboundMessage(message, key);
const text = promptMessage.content;
const hasImages = hasInboundImages(promptMessage);
const hasFiles = hasInboundFiles(promptMessage);
const hasReply = hasReplyReference(promptMessage);
const isPlainText = String(message?.msgtype).toLowerCase() === 'text';
let cardStream = null;
let cardStarted = false;
let cardStartedAt = null;
let batchSettled = batchSubmission === null;
try {
if (!text && !hasImages && !hasFiles) {
if (!text && !hasImages && !hasFiles && !hasReply) {
await this.#send(sessionWebhook, t('目前支持文字、图片和文件消息。'), this.#atUsersFor(message));
return;
}
@ -992,8 +1123,8 @@ export class DingtalkHarnessBridge {
return;
}
let content = hasImages
? await promptContentForMessage(promptMessage, { signal: this.#signal })
let content = hasImages || hasReply
? await promptContentForInboundMessage(promptMessage, { signal: this.#signal })
: undefined;
const snapshot = this.#acceptedMessageIds.get(messageId);
let contextEnhanced = false;
@ -1018,7 +1149,9 @@ export class DingtalkHarnessBridge {
signal: this.#signal,
logger: this.#logger,
});
const startedAt = Date.now();
cardStarted = await cardStream.start(t(CARD_INITIAL_TEXT));
if (cardStarted) cardStartedAt = startedAt;
}
const { answer, artifacts = [] } = await askInWorkspaceSession({
harness: this.#harness,
@ -1056,12 +1189,14 @@ export class DingtalkHarnessBridge {
let textDeliveryError = null;
let textReceipt = null;
let streamed = false;
const deliveryStartedAt = cardStartedAt ?? Date.now();
try {
streamed = cardStarted && await cardStream.finish(answerText);
if (streamed) {
textReceipt = createDeliveryReceipt({
deliveryId: messageId,
presentation: 'dingtalk-card',
providerMessageIds: cardStream.providerMessageIds,
});
} else {
textReceipt = createDeliveryReceipt({
@ -1070,6 +1205,17 @@ export class DingtalkHarnessBridge {
providerMessageIds: await this.#send(sessionWebhook, answerText, this.#atUsersFor(message)),
});
}
try {
await this.#state.rememberOutboundMessage?.({
conversationKey: key,
text: answerText,
sentAt: deliveryStartedAt,
completedAt: Date.now(),
providerMessageIds: providerMessageIdsFor(textReceipt),
});
} catch (error) {
this.#logger.warn?.('[dsh-dingtalk] failed to remember an outbound message:', error);
}
} catch (error) {
textDeliveryError = channelDeliveryFailure(error);
}

View file

@ -32,7 +32,7 @@ function requiredCredential(value, name) {
* @param {number} [options.updateIntervalMs=500] Minimum delay between updates.
* @param {()=>number} [options.clock] Monotonic millisecond clock.
* @param {{setTimeout: Function, clearTimeout: Function}} [options.timer] Timer implementation.
* @returns {{start(initialText: string): Promise<boolean>, push(progressText: string): void, finish(finalText: string): Promise<boolean>}}
* @returns {{start(initialText: string): Promise<boolean>, push(progressText: string): void, finish(finalText: string): Promise<boolean>, readonly providerMessageIds: string[]}}
* Card stream controller.
*/
export function createDingTalkCardStream({
@ -231,5 +231,12 @@ export function createDingTalkCardStream({
return finishPromise;
};
return Object.freeze({ start, push, finish });
return Object.freeze({
start,
push,
finish,
get providerMessageIds() {
return cardRequest?.cardInstanceId ? [cardRequest.cardInstanceId] : [];
},
});
}

View file

@ -9,8 +9,14 @@ const EMPTY_STATE = Object.freeze({
sessions: {},
seenMessageIds: [],
pendingSenders: {},
recentOutboundMessages: [],
});
export const DINGTALK_RECENT_OUTBOUND_LIMIT = 200;
export const DINGTALK_RECENT_OUTBOUND_TTL_MS = 30 * 24 * 60 * 60 * 1_000;
export const DINGTALK_RECENT_OUTBOUND_TEXT_LIMIT = 8_000;
export const DINGTALK_RECENT_OUTBOUND_MATCH_TOLERANCE_MS = 15_000;
function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
@ -19,6 +25,49 @@ function displayName(value) {
return (nonEmptyString(value) ?? t('钉钉用户')).slice(0, 100);
}
function timestampMs(value) {
const number = typeof value === 'string' && value.trim() ? Number(value) : value;
if (!Number.isFinite(number) || number < 0) return null;
return Math.trunc(number < 10_000_000_000 ? number * 1_000 : number);
}
function providerMessageIds(value) {
if (!Array.isArray(value)) return [];
return [...new Set(value
.map((id) => (id === undefined || id === null ? null : nonEmptyString(String(id))))
.filter(Boolean))];
}
function truncateText(value) {
const text = nonEmptyString(value);
return text ? [...text].slice(0, DINGTALK_RECENT_OUTBOUND_TEXT_LIMIT).join('') : null;
}
function normalizeRecentOutboundMessage(value) {
if (!value || typeof value !== 'object' || Array.isArray(value)) return null;
const conversationKey = nonEmptyString(value.conversationKey);
const text = truncateText(value.text);
const sentAt = timestampMs(value.sentAt);
const completedAt = timestampMs(value.completedAt) ?? sentAt;
if (!conversationKey || !text || sentAt === null || completedAt === null) return null;
return {
conversationKey,
text,
sentAt,
completedAt: Math.max(sentAt, completedAt),
providerMessageIds: providerMessageIds(value.providerMessageIds),
};
}
function recentOutboundMessages(value, now = Date.now()) {
if (!Array.isArray(value)) return [];
const cutoff = now - DINGTALK_RECENT_OUTBOUND_TTL_MS;
return value
.map(normalizeRecentOutboundMessage)
.filter((entry) => entry && entry.completedAt >= cutoff)
.slice(-DINGTALK_RECENT_OUTBOUND_LIMIT);
}
function normalizePendingSender(value, fallbackRequestId) {
if (!value || typeof value !== 'object') return null;
const requestId = nonEmptyString(value.requestId) ?? nonEmptyString(fallbackRequestId);
@ -69,6 +118,7 @@ function normalizeState(value) {
? [...new Set(value.seenMessageIds.map(nonEmptyString).filter(Boolean))].slice(-1_000)
: [],
pendingSenders,
recentOutboundMessages: recentOutboundMessages(value.recentOutboundMessages),
};
}
@ -140,6 +190,54 @@ export class DingtalkStateStore {
await this.#persist();
}
async rememberOutboundMessage({
conversationKey,
text,
sentAt = Date.now(),
completedAt = Date.now(),
providerMessageIds: messageIds = [],
} = {}) {
const entry = normalizeRecentOutboundMessage({
conversationKey,
text,
sentAt,
completedAt,
providerMessageIds: messageIds,
});
if (!entry) throw new TypeError('Invalid DingTalk outbound message');
this.#state.recentOutboundMessages = recentOutboundMessages([
...this.#state.recentOutboundMessages,
entry,
]);
await this.#persist();
}
recentOutboundTextFor({
conversationKey,
processQueryKey,
messageId,
createdAt,
now = Date.now(),
} = {}) {
const key = nonEmptyString(conversationKey);
if (!key) return null;
const active = recentOutboundMessages(this.#state.recentOutboundMessages, now)
.filter((entry) => entry.conversationKey === key);
const quotedIds = providerMessageIds([processQueryKey, messageId]);
for (const quotedId of quotedIds) {
const exact = active.filter((entry) => entry.providerMessageIds.includes(quotedId));
if (exact.length === 1) return exact[0].text;
if (exact.length > 1) return null;
}
const quotedAt = timestampMs(createdAt);
if (quotedAt === null) return null;
const candidates = active.filter((entry) => (
quotedAt >= entry.sentAt - DINGTALK_RECENT_OUTBOUND_MATCH_TOLERANCE_MS
&& quotedAt <= entry.completedAt + DINGTALK_RECENT_OUTBOUND_MATCH_TOLERANCE_MS
));
return candidates.length === 1 ? candidates[0].text : null;
}
pendingSenders() {
return Object.values(this.#state.pendingSenders)
.sort((left, right) => left.requestedAt.localeCompare(right.requestedAt))