feat: add status reactions across messaging channels

This commit is contained in:
xmanrui 2026-08-25 17:12:28 +08:00
parent aa8fd71b93
commit da01e3acee
30 changed files with 1853 additions and 274 deletions

View file

@ -8,8 +8,14 @@ export const DINGTALK_REGISTRATION_BASE_URL = 'https://oapi.dingtalk.com/';
export const DINGTALK_API_BASE_URL = 'https://api.dingtalk.com/';
export const DINGTALK_REGISTRATION_SOURCE = 'DING_DWS_CLAW';
export const DINGTALK_AI_CARD_TEMPLATE_ID = '02fcf2f4-5e02-4a85-b672-46d1f715543e.schema';
export const DINGTALK_THINKING_REACTION_NAME = '🤔思考中';
export const DINGTALK_DONE_REACTION_NAME = '✅已完成';
export const DINGTALK_ERROR_REACTION_NAME = '❌处理失败';
const DEFAULT_TIMEOUT_MS = 15_000;
const REACTION_TIMEOUT_MS = 5_000;
const TEXT_REACTION_ID = '2659900';
const TEXT_REACTION_BACKGROUND_ID = 'im_bg_1';
const REGISTRATION_STATUSES = new Set(['WAITING', 'SUCCESS', 'FAIL', 'EXPIRED']);
export class DingtalkApiError extends Error {
@ -406,18 +412,20 @@ export function createDingtalkApi({
if (!source) throw new TypeError('registrationSource is required');
const tokenCache = new Map();
const tokenRequests = new Map();
const reactionTokenRequests = new Map();
let cardSlotTail = Promise.resolve();
let nextCardRequestAt = 0;
const endpoint = (base, pathname) => new URL(pathname.replace(/^\//, ''), base);
async function accessToken({ clientId, clientSecret, signal }) {
async function accessToken({ clientId, clientSecret, signal, requestKind = 'normal' }) {
const appKey = nonEmptyString(clientId);
const appSecret = nonEmptyString(clientSecret);
if (!appKey || !appSecret) throw new TypeError('clientId and clientSecret are required');
const cached = tokenCache.get(appKey);
if (cached && cached.expiresAt > now()) return cached.token;
if (tokenRequests.has(appKey)) return tokenRequests.get(appKey);
const requests = requestKind === 'reaction' ? reactionTokenRequests : tokenRequests;
if (requests.has(appKey)) return requests.get(appKey);
const request = (async () => {
const value = await requestJson(fetchImpl, endpoint(apiBase, 'v1.0/oauth2/accessToken'), {
@ -431,9 +439,75 @@ export function createDingtalkApi({
const refreshAfterMs = Math.max(1_000, (expiresInSeconds - 60) * 1_000);
tokenCache.set(appKey, { token, expiresAt: now() + refreshAfterMs });
return token;
})().finally(() => tokenRequests.delete(appKey));
tokenRequests.set(appKey, request);
return request;
})();
const shared = request.finally(() => requests.delete(appKey));
requests.set(appKey, shared);
return shared;
}
async function changeReaction({
clientId,
clientSecret,
robotCode,
messageId,
conversationId,
reactionName = DINGTALK_THINKING_REACTION_NAME,
signal,
}, action) {
const appKey = nonEmptyString(clientId);
const appSecret = nonEmptyString(clientSecret);
const botCode = nonEmptyString(robotCode) ?? appKey;
const openMsgId = nonEmptyString(messageId);
const openConversationId = nonEmptyString(conversationId);
const emotionName = nonEmptyString(reactionName);
if (!appKey || !appSecret) throw new TypeError('clientId and clientSecret are required');
if (!botCode) throw new TypeError('robotCode is required');
if (!openMsgId || !openConversationId) {
throw new TypeError('messageId and conversationId are required');
}
if (!emotionName) throw new TypeError('reactionName is required');
// Reactions deduplicate with each other but never own the normal reply's
// in-flight token request. Both paths may still reuse a cached token.
const token = await accessToken({
clientId: appKey,
clientSecret: appSecret,
signal,
requestKind: 'reaction',
});
const response = await requestJson(
fetchImpl,
endpoint(apiBase, `v1.0/robot/emotion/${action}`),
{
body: {
robotCode: botCode,
openMsgId,
openConversationId,
emotionType: 2,
emotionName,
textEmotion: {
emotionId: TEXT_REACTION_ID,
emotionName,
text: emotionName,
backgroundId: TEXT_REACTION_BACKGROUND_ID,
},
},
headers: { 'x-acs-dingtalk-access-token': token },
signal,
timeoutMs: REACTION_TIMEOUT_MS,
action: action === 'reply' ? '消息状态添加' : '消息状态撤回',
},
);
const rejection = response?.success === false
? safeProviderCode(response?.code ?? response?.errcode) ?? 'rejected'
: rejectedProviderResponse(response);
if (rejection) {
throw new DingtalkApiError(
'reaction-rejected',
'钉钉服务拒绝了消息状态请求。',
{ providerCode: rejection },
);
}
return true;
}
async function messageFileDownloadUrl({
@ -690,6 +764,28 @@ export function createDingtalkApi({
accessToken,
async addReaction(request) {
return changeReaction(request, 'reply');
},
async recallReaction(request) {
return changeReaction(request, 'recall');
},
async addThinkingReaction(request) {
return changeReaction({
...request,
reactionName: DINGTALK_THINKING_REACTION_NAME,
}, 'reply');
},
async recallThinkingReaction(request) {
return changeReaction({
...request,
reactionName: DINGTALK_THINKING_REACTION_NAME,
}, 'recall');
},
async downloadImage({
clientId,
clientSecret,

View file

@ -1,4 +1,7 @@
import {
DINGTALK_DONE_REACTION_NAME,
DINGTALK_ERROR_REACTION_NAME,
DINGTALK_THINKING_REACTION_NAME,
normalizeDingtalkSessionWebhook,
splitDingtalkText,
} from './dingtalk-api.mjs';
@ -309,7 +312,15 @@ function canClaimInteractionReply(message, pending, sender) {
function ensureStats(status) {
status.stats ??= {};
for (const key of ['messagesReceived', 'messagesReplied', 'messagesRejected', 'messagesIgnored']) {
for (const key of [
'messagesReceived',
'messagesReplied',
'messagesRejected',
'messagesIgnored',
'reactionsAdded',
'reactionsRemoved',
'reactionErrors',
]) {
status[key] ??= 0;
status.stats[key] = status[key];
}
@ -328,6 +339,9 @@ export function createDingtalkBridgeStatus({ pendingSenders = [] } = {}) {
messagesReplied: 0,
messagesRejected: 0,
messagesIgnored: 0,
reactionsAdded: 0,
reactionsRemoved: 0,
reactionErrors: 0,
lastMessageAt: null,
lastReplyAt: null,
lastRejectedAt: null,
@ -339,6 +353,9 @@ export function createDingtalkBridgeStatus({ pendingSenders = [] } = {}) {
messagesReplied: 0,
messagesRejected: 0,
messagesIgnored: 0,
reactionsAdded: 0,
reactionsRemoved: 0,
reactionErrors: 0,
},
};
}
@ -352,6 +369,7 @@ export class DingtalkHarnessBridge {
#status;
#logger;
#replyTimeoutMs;
#reactionTimeoutMs;
#maxMessageChars;
#signal;
#queues = new Map();
@ -372,6 +390,7 @@ export class DingtalkHarnessBridge {
status = createDingtalkBridgeStatus(),
logger = console,
replyTimeoutMs = 600_000,
reactionTimeoutMs = 5_000,
maxMessageChars = 4_000,
signal,
}) {
@ -389,6 +408,9 @@ export class DingtalkHarnessBridge {
this.#logger = logger;
this.#approvals = new HarnessApprovalQueue({ label: 'DingTalk', logger });
this.#replyTimeoutMs = replyTimeoutMs;
this.#reactionTimeoutMs = Number.isFinite(reactionTimeoutMs) && reactionTimeoutMs > 0
? Math.floor(reactionTimeoutMs)
: 5_000;
this.#maxMessageChars = maxMessageChars;
this.#signal = signal;
ensureStats(this.#status);
@ -436,15 +458,35 @@ export class DingtalkHarnessBridge {
const commandText = nonEmptyString(promptMessage.content) ?? '';
const addressed = String(message.conversationType) !== '2' || message?.isInAtList === true;
const direct = String(message.conversationType) !== '2';
const statusReaction = sessionWebhook && addressed ? this.#startStatusReaction(message) : null;
const finish = (task) => Promise.resolve(task).then(
(value) => {
this.#finishStatusReaction(
statusReaction,
this.#signal?.aborted ? 'clear' : 'success',
);
return value;
},
(error) => {
this.#finishStatusReaction(
statusReaction,
this.#signal?.aborted || error?.name === 'AbortError' || error?.code === 'turn-stopped'
? 'clear'
: 'error',
);
throw error;
},
);
const batchCommand = String(message?.msgtype).toLowerCase() === 'text'
&& isBatchInputCommand(commandText);
const batchStatus = this.#batchInputs.status(key);
if (batchCommand && !direct && sessionWebhook && addressed) {
return this.#finishBatchResult(
return finish(this.#finishBatchResult(
messageId,
sessionWebhook,
{ message: batchInputGroupUnsupportedMessage() },
);
statusReaction,
));
}
if (direct && sessionWebhook && (batchCommand || batchStatus.phase === 'collecting')) {
const exactBatchStart = /^\/batch$/iu.test(commandText);
@ -460,7 +502,7 @@ export class DingtalkHarnessBridge {
});
if (result.handled) {
if (result.kind === 'submit') {
return this.#enqueueMessage(
return finish(this.#enqueueMessage(
{
...message,
msgtype: 'text',
@ -469,10 +511,15 @@ export class DingtalkHarnessBridge {
messageId,
sender,
key,
{ batchSubmission: result },
);
{ batchSubmission: result, statusReaction },
));
}
return this.#finishBatchResult(messageId, sessionWebhook, result);
return finish(this.#finishBatchResult(
messageId,
sessionWebhook,
result,
statusReaction,
));
}
}
const commandRunner = hasInboundFiles(promptMessage) ? null : isControlCommand(commandText)
@ -490,7 +537,11 @@ export class DingtalkHarnessBridge {
promptMessage,
commandRunner,
).catch((error) => {
if (error?.code === 'turn-stopped' || this.#signal?.aborted) return;
if (error?.code === 'turn-stopped' || this.#signal?.aborted) {
this.#finishStatusReaction(statusReaction, 'clear');
return;
}
this.#finishStatusReaction(statusReaction, 'error');
this.#status.lastError = error?.message ?? String(error);
const failure = setLastMessageFailure(this.#status, error);
this.#logger.error?.(
@ -503,7 +554,7 @@ export class DingtalkHarnessBridge {
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
return finish(task);
}
const approvalReply = this.#approvals.claimReply({
key,
@ -537,7 +588,11 @@ export class DingtalkHarnessBridge {
return true;
})
.catch((error) => {
if (this.#signal?.aborted) return;
if (this.#signal?.aborted) {
this.#finishStatusReaction(statusReaction, 'clear');
return;
}
this.#finishStatusReaction(statusReaction, 'error');
this.#status.lastError = t('钉钉审批处理失败。');
this.#logger.error?.('[dsh-dingtalk] failed to process an approval reply', error);
})
@ -546,18 +601,18 @@ export class DingtalkHarnessBridge {
this.#interactionTasks.delete(current);
});
this.#interactionTasks.add(current);
return current;
return finish(current);
}
if (pending && pending.actor !== sender) {
return this.#enqueueMessage(message, messageId, sender, key);
return finish(this.#enqueueMessage(message, messageId, sender, key, { statusReaction }));
}
// Once one valid answer has been claimed, later messages are subsequent
// prompts even if the network submission eventually needs a retry. Invalid
// replies do not claim the question, so the next valid answer can still
// pass through this interaction queue.
if (pending?.submitting || pending?.claimedReplyMessageId) {
return this.#enqueueMessage(message, messageId, sender, key);
return finish(this.#enqueueMessage(message, messageId, sender, key, { statusReaction }));
}
if (pending) {
if (canClaimInteractionReply(message, pending, sender)) {
@ -566,7 +621,14 @@ export class DingtalkHarnessBridge {
const previous = pending.queue ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#processInteractionReply(message, messageId, sender, key, pending))
.then(() => this.#processInteractionReply(
message,
messageId,
sender,
key,
pending,
statusReaction,
))
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
if (pending.claimedReplyMessageId === messageId) {
@ -575,15 +637,101 @@ export class DingtalkHarnessBridge {
if (pending.queue === current) pending.queue = null;
});
pending.queue = current;
return current;
return finish(current);
}
return this.#enqueueMessage(message, messageId, sender, key);
return finish(this.#enqueueMessage(message, messageId, sender, key, { statusReaction }));
}
#runReactionCall(method, target, reactionName, kind) {
const controller = new AbortController();
const operation = Promise.resolve().then(() => this.#api[method]({
clientId: this.#clientId,
clientSecret: this.#clientSecret,
...target,
reactionName,
signal: controller.signal,
}));
let timer;
const timeout = new Promise((_, reject) => {
timer = setTimeout(() => {
const error = new DOMException('DingTalk reaction timed out', 'TimeoutError');
controller.abort(error);
reject(error);
}, this.#reactionTimeoutMs);
timer.unref?.();
});
return Promise.race([operation, timeout])
.then(() => {
increment(this.#status, kind === 'add' ? 'reactionsAdded' : 'reactionsRemoved');
return true;
})
.catch((error) => {
increment(this.#status, 'reactionErrors');
this.#logger.debug?.(`[dsh-dingtalk] ${method} failed`, safeErrorDiagnostic(error));
return false;
})
.finally(() => clearTimeout(timer));
}
#startStatusReaction(message) {
if (typeof this.#api.addReaction !== 'function'
|| typeof this.#api.recallReaction !== 'function') return null;
const messageId = nonEmptyString(message?.msgId);
const conversationId = nonEmptyString(message?.conversationId);
if (!messageId || !conversationId) return null;
const target = {
messageId,
conversationId,
robotCode: nonEmptyString(message?.robotCode) ?? this.#clientId,
};
return {
target,
attached: this.#runReactionCall(
'addReaction',
target,
DINGTALK_THINKING_REACTION_NAME,
'add',
),
terminal: false,
};
}
#finishStatusReaction(reaction, outcome) {
if (!reaction || reaction.terminal) return;
reaction.terminal = true;
const terminalName = outcome === 'success'
? DINGTALK_DONE_REACTION_NAME
: outcome === 'error' ? DINGTALK_ERROR_REACTION_NAME : null;
// Preserve attach -> recall -> terminal ordering without extending the message task.
void reaction.attached.then(async (attached) => {
let cleaned = await this.#runReactionCall(
'recallReaction',
reaction.target,
DINGTALK_THINKING_REACTION_NAME,
'remove',
);
if (!attached || !cleaned) {
await new Promise((resolve) => {
const retry = setTimeout(resolve, Math.min(1_000, this.#reactionTimeoutMs));
retry.unref?.();
});
cleaned = await this.#runReactionCall(
'recallReaction',
reaction.target,
DINGTALK_THINKING_REACTION_NAME,
'remove',
) || cleaned;
}
if (!cleaned || !terminalName || this.#signal?.aborted) return;
await this.#runReactionCall('addReaction', reaction.target, terminalName, 'add');
}).catch(() => undefined);
}
#enqueueMessage(message, messageId, sender, key, {
releaseMessageId = true,
alreadyRecorded = false,
batchSubmission = null,
statusReaction = null,
} = {}) {
let hasSafeReplyRoute = false;
try {
@ -607,6 +755,7 @@ export class DingtalkHarnessBridge {
alreadyRecorded,
preparedMessage,
batchSubmission,
statusReaction,
}))
.finally(() => {
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
@ -659,7 +808,7 @@ export class DingtalkHarnessBridge {
this.#status.lastError = null;
}
#finishBatchResult(messageId, sessionWebhook, result) {
#finishBatchResult(messageId, sessionWebhook, result, statusReaction) {
let task;
task = Promise.resolve().then(async () => {
if (this.#state.hasSeen(messageId)) return;
@ -669,7 +818,11 @@ export class DingtalkHarnessBridge {
if (result.message) await this.#send(sessionWebhook, result.message);
this.#status.lastError = null;
}).catch(async (error) => {
if (this.#signal?.aborted) return;
if (this.#signal?.aborted) {
this.#finishStatusReaction(statusReaction, 'clear');
return;
}
this.#finishStatusReaction(statusReaction, 'error');
this.#status.lastError = error?.message ?? String(error);
const failure = setLastMessageFailure(this.#status, error);
this.#logger.error?.(
@ -689,6 +842,7 @@ export class DingtalkHarnessBridge {
alreadyRecorded = false,
preparedMessage,
batchSubmission = null,
statusReaction = null,
} = {}) {
this.#signal?.throwIfAborted();
if (!alreadyRecorded) {
@ -865,10 +1019,15 @@ export class DingtalkHarnessBridge {
batchSettled = true;
}
if (error?.code === 'turn-stopped') {
this.#finishStatusReaction(statusReaction, 'clear');
if (cardStarted) await cardStream.finish(t('已停止。')).catch(() => undefined);
return;
}
if (this.#signal?.aborted) return;
if (this.#signal?.aborted) {
this.#finishStatusReaction(statusReaction, 'clear');
return;
}
this.#finishStatusReaction(statusReaction, 'error');
this.#status.lastError = error?.message ?? String(error);
const userMessage = inboundFileUserMessage(error)
?? dingtalkImageErrorUserMessage(error);
@ -896,7 +1055,14 @@ export class DingtalkHarnessBridge {
}
}
async #processInteractionReply(message, messageId, sender, key, expected) {
async #processInteractionReply(
message,
messageId,
sender,
key,
expected,
statusReaction,
) {
this.#signal?.throwIfAborted();
const current = this.#pendingInteractions.get(key);
const claimed = expected.claimedReplyMessageId === messageId;
@ -904,7 +1070,10 @@ export class DingtalkHarnessBridge {
if (claimed && (!current || current !== expected)) {
return this.#discardResolvedInteractionReply(message, messageId);
}
return this.#enqueueMessage(message, messageId, sender, key, { releaseMessageId: false });
return this.#enqueueMessage(message, messageId, sender, key, {
releaseMessageId: false,
statusReaction,
});
}
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
@ -949,6 +1118,7 @@ export class DingtalkHarnessBridge {
return this.#enqueueMessage(message, messageId, sender, key, {
releaseMessageId: false,
alreadyRecorded: true,
statusReaction,
});
}
pending.sessionWebhook = sessionWebhook;
@ -956,6 +1126,7 @@ export class DingtalkHarnessBridge {
try {
await this.#presentInteraction(pending);
} catch {
this.#finishStatusReaction(statusReaction, 'error');
this.#status.lastError = t('钉钉交互问题发送失败。');
this.#logger.error?.('[dsh-dingtalk] failed to retry an interaction question');
pending.interaction.reconnect?.();
@ -975,6 +1146,7 @@ export class DingtalkHarnessBridge {
try {
await this.#presentInteraction(pending);
} catch {
this.#finishStatusReaction(statusReaction, 'error');
this.#status.lastError = t('钉钉交互问题发送失败。');
this.#logger.error?.('[dsh-dingtalk] failed to send the next interaction question');
pending.interaction.reconnect?.();
@ -994,7 +1166,10 @@ export class DingtalkHarnessBridge {
this.#clearPendingInteraction(key, pending.interactionId);
this.#status.lastError = null;
} catch (error) {
if (this.#signal?.aborted) return;
if (this.#signal?.aborted) {
this.#finishStatusReaction(statusReaction, 'clear');
return;
}
if (this.#pendingInteractions.get(key) !== pending) return;
if (error?.code === 'interaction-not-pending') {
this.#clearPendingInteraction(key, pending.interactionId);
@ -1005,6 +1180,7 @@ export class DingtalkHarnessBridge {
}
return;
}
this.#finishStatusReaction(statusReaction, 'error');
pending.submitting = false;
pending.answers.pop();
pending.index -= 1;

View file

@ -88,6 +88,12 @@ function snowflake(value, name) {
return id;
}
function reactionEmoji(value) {
const emoji = cleanString(value);
if (!emoji) throw new TypeError('A Discord reaction emoji is required');
return emoji;
}
export function validDiscordToken(value) {
return typeof value === 'string'
&& /^[A-Za-z0-9_-]{20,}\.[A-Za-z0-9_-]{4,}\.[A-Za-z0-9_-]{20,}$/.test(value.trim());
@ -229,6 +235,35 @@ export class DiscordApi {
});
}
async addOwnReaction({ channelId, messageId, emoji, signal } = {}) {
const normalizedEmoji = reactionEmoji(emoji);
await this.#request(
`channels/${snowflake(channelId, 'channel id')}/messages/${snowflake(messageId, 'message id')}`
+ `/reactions/${encodeURIComponent(normalizedEmoji)}/@me`,
{
method: 'PUT',
signal,
expectBody: false,
retry: false,
},
);
return normalizedEmoji;
}
removeOwnReaction({ channelId, messageId, emoji, signal } = {}) {
const normalizedEmoji = reactionEmoji(emoji);
return this.#request(
`channels/${snowflake(channelId, 'channel id')}/messages/${snowflake(messageId, 'message id')}`
+ `/reactions/${encodeURIComponent(normalizedEmoji)}/@me`,
{
method: 'DELETE',
signal,
expectBody: false,
retry: false,
},
);
}
async #request(path, {
method,
body,

View file

@ -4,6 +4,7 @@ export const DISCORD_DESCRIPTOR = Object.freeze({
key: 'discord',
label: 'Discord',
connectionLabel: ' Gateway 长连接',
reactions: Object.freeze({ processing: '👀', success: '✅', error: '❌' }),
});
export class DiscordHarnessBridge extends TextHarnessBridge {

View file

@ -231,6 +231,10 @@ export function normalizeDiscordMessage(message, botId, { fetchImpl = fetch } =
channelId: String(message.channel_id),
replyToMessageId: String(message.id),
},
reactionTarget: {
channelId: String(message.channel_id),
messageId: String(message.id),
},
connectionTestTarget: { channelId: String(message.channel_id) },
};
}
@ -350,6 +354,24 @@ export class DiscordBotClient {
});
}
addReaction(target, emoji, { signal } = {}) {
return this.#api.addOwnReaction({
channelId: target.channelId,
messageId: target.messageId,
emoji,
signal: this.#operationSignal(signal),
});
}
removeReaction(target, reactionKey, { signal } = {}) {
return this.#api.removeOwnReaction({
channelId: target.channelId,
messageId: target.messageId,
emoji: reactionKey,
signal: this.#operationSignal(signal),
});
}
async openStream(target) {
const notice = !this.#deliveredNotices.has(target) && target?.notice
? String(target.notice) : null;
@ -381,6 +403,10 @@ export class DiscordBotClient {
});
return stream.start();
}
#operationSignal(signal) {
return signal ?? this.#signal;
}
}
export function createDiscordRuntimeStatus() {

View file

@ -55,6 +55,7 @@ import {
messageFailureText,
setLastMessageFailure,
} from '../shared/message-failure.mjs';
import { beginStatusReaction } from '../shared/status-reaction.mjs';
import {
MENU_PAGE_SIZE,
PRESET_FOLLOW_DEFAULT_SENTINEL,
@ -549,7 +550,7 @@ export class FeishuHarnessBridge {
}
this.#acceptedMessageIds.add(messageId);
const processingReaction = this.#addReaction(messageId, 'OnIt');
const processingReaction = this.#beginReaction(messageId);
const commandMessage = extractInboundMessage(event, this.#client);
const commandText = nonEmptyString(commandMessage.content) ?? '';
const batchText = event.message.message_type === 'text'
@ -3603,35 +3604,24 @@ export class FeishuHarnessBridge {
return t('_{text}_', { text: update.text || t('正在处理…') });
}
async #addReaction(messageId, emojiType) {
if (!this.#channel?.addReaction) return null;
try {
const reactionId = await this.#channel.addReaction(messageId, emojiType);
this.#status.reactionsAdded = (this.#status.reactionsAdded ?? 0) + 1;
return reactionId;
} catch (error) {
this.#status.reactionErrors = (this.#status.reactionErrors ?? 0) + 1;
this.#logger.warn?.(`[dsh-feishu] unable to add ${emojiType} reaction:`, error.message);
return null;
}
#beginReaction(messageId) {
return beginStatusReaction({
adapter: this.#channel,
target: messageId,
reactions: { processing: 'OnIt', success: 'DONE', error: 'ERROR' },
status: this.#status,
logger: this.#logger,
label: 'feishu',
});
}
async #removeProcessingReaction(messageId, processingReaction) {
const reactionId = await processingReaction;
if (reactionId && this.#channel?.removeReaction) {
try {
await this.#channel.removeReaction(messageId, reactionId);
this.#status.reactionsRemoved = (this.#status.reactionsRemoved ?? 0) + 1;
} catch (error) {
this.#status.reactionErrors = (this.#status.reactionErrors ?? 0) + 1;
this.#logger.warn?.('[dsh-feishu] unable to remove processing reaction:', error.message);
}
}
#removeProcessingReaction(_messageId, processingReaction) {
processingReaction.clear();
}
async #finishReaction(messageId, processingReaction, finalEmojiType) {
await this.#removeProcessingReaction(messageId, processingReaction);
await this.#addReaction(messageId, finalEmojiType);
#finishReaction(_messageId, processingReaction, finalEmojiType) {
if (finalEmojiType === 'ERROR') processingReaction.error();
else processingReaction.success();
}
async #send(chatId, text) {

View file

@ -0,0 +1,107 @@
const DEFAULT_TIMEOUT_MS = 2_000;
const NOOP_REACTION = Object.freeze({
success() {},
error() {},
clear() {},
settled: () => Promise.resolve(),
});
function increment(status, key) {
if (!status || typeof status !== 'object') return;
status[key] = (status[key] ?? 0) + 1;
}
async function runWithTimeout(operation, timeoutMs) {
const signal = AbortSignal.timeout(timeoutMs);
let onAbort;
const aborted = new Promise((_, reject) => {
onAbort = () => reject(signal.reason ?? new DOMException('Timed out', 'TimeoutError'));
signal.addEventListener('abort', onAbort, { once: true });
});
try {
return await Promise.race([operation(signal), aborted]);
} finally {
signal.removeEventListener('abort', onAbort);
}
}
/**
* Starts a best-effort status reaction lifecycle which is deliberately not
* part of the caller's message queue. Calls are serialized only for this one
* source message; every provider operation is bounded and absorbs failures.
*/
export function beginStatusReaction({
adapter,
target,
reactions,
status,
logger = console,
label = 'channel',
timeoutMs = DEFAULT_TIMEOUT_MS,
} = {}) {
if (!target
|| typeof adapter?.addReaction !== 'function'
|| typeof adapter?.removeReaction !== 'function'
|| typeof reactions?.processing !== 'string'
|| !reactions.processing
|| !Number.isSafeInteger(timeoutMs)
|| timeoutMs <= 0) return NOOP_REACTION;
let currentReaction = null;
let terminal = false;
const safely = async (kind, operation) => {
try {
const value = await runWithTimeout(operation, timeoutMs);
increment(status, kind === 'add' ? 'reactionsAdded' : 'reactionsRemoved');
return { ok: true, value };
} catch (cause) {
increment(status, 'reactionErrors');
logger.warn?.(
`[dsh-im:${label}] status reaction ${kind} failed:`,
cause?.message ?? cause?.name ?? String(cause),
);
return { ok: false, value: null };
}
};
const transition = async (emoji) => {
if (currentReaction !== null) {
const previous = currentReaction;
currentReaction = null;
await safely('remove', (signal) => adapter.removeReaction(
target,
previous,
{ signal },
));
}
if (typeof emoji !== 'string' || !emoji) return;
const added = await safely('add', (signal) => adapter.addReaction(
target,
emoji,
{ signal },
));
if (added.ok && added.value !== undefined && added.value !== null) {
currentReaction = added.value;
}
};
// Calling an async function starts the provider request synchronously up to
// its first await, while the returned tail remains completely detached from
// normal message processing.
let tail = transition(reactions.processing);
const finish = (emoji) => {
if (terminal) return;
terminal = true;
tail = tail.then(() => transition(emoji), () => transition(emoji));
void tail.catch(() => undefined);
};
return Object.freeze({
success: () => finish(reactions.success),
error: () => finish(reactions.error),
clear: () => finish(null),
settled: () => tail,
});
}

View file

@ -52,6 +52,7 @@ import {
messageFailureText,
setLastMessageFailure,
} from './message-failure.mjs';
import { beginStatusReaction } from './status-reaction.mjs';
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
const FILE_ONLY_COMPLETION_TEXT = '任务已完成。';
@ -114,6 +115,9 @@ export function createTextBridgeStatus() {
lastRejectedAt: null,
lastError: null,
lastMessageError: null,
reactionsAdded: 0,
reactionsRemoved: 0,
reactionErrors: 0,
};
}
@ -178,6 +182,27 @@ export class TextHarnessBridge {
return Promise.resolve();
}
this.#acceptedMessageIds.add(messageId);
const statusReaction = beginStatusReaction({
adapter: this.#bot,
target: normalized.kind === 'direct' || normalized.addressed === true
? normalized.reactionTarget
: null,
reactions: this.#descriptor.reactions,
status: this.#status,
logger: this.#logger,
label: this.#descriptor.key,
});
normalized.statusReaction = statusReaction;
const processing = this.#acceptAcceptedMessage(normalized, messageId, senderId);
void processing.then(
() => statusReaction.success(),
() => statusReaction.error(),
);
return processing;
}
#acceptAcceptedMessage(normalized, messageId, senderId) {
if (normalized.kind === 'direct') {
rememberConnectionTestTarget(
this.#state,
@ -185,7 +210,7 @@ export class TextHarnessBridge {
);
}
const key = `${kind}:${conversationId}`;
const key = `${normalized.kind}:${normalized.conversationId}`;
const pending = this.#pendingInteractions.get(key);
const text = cleanText(normalized.content);
const batchCommand = isBatchInputCommand(text);
@ -326,7 +351,11 @@ export class TextHarnessBridge {
if (reply) await this.#bot.sendText(message.replyTarget, reply);
this.#status.lastError = null;
})().catch(async (error) => {
if (this.#signal?.aborted) return;
if (this.#signal?.aborted) {
message.statusReaction?.clear();
return;
}
message.statusReaction?.error();
this.#status.lastError = error?.message ?? String(error);
const failure = setLastMessageFailure(this.#status, error);
this.#logger.error?.(
@ -406,7 +435,11 @@ export class TextHarnessBridge {
}
this.#status.lastError = null;
} catch (error) {
if (error?.code === 'turn-stopped' || this.#signal?.aborted) return;
if (error?.code === 'turn-stopped' || this.#signal?.aborted) {
message.statusReaction?.clear();
return;
}
message.statusReaction?.error();
this.#status.lastError = error?.message ?? String(error);
const failure = setLastMessageFailure(this.#status, error);
this.#logger.error?.(
@ -706,6 +739,7 @@ export class TextHarnessBridge {
? this.#batches.fail(conversationKey, batchSubmission.token)
: null;
if (turnStopped) {
message.statusReaction?.clear();
if (stream) {
try {
await stream.finish(t('已停止。'));
@ -716,9 +750,11 @@ export class TextHarnessBridge {
return;
}
if (this.#signal?.aborted) {
message.statusReaction?.clear();
stream?.cancel?.();
return;
}
message.statusReaction?.error();
this.#status.lastError = error?.message ?? String(error);
const presentStreamFailure = async (text) => {
const method = typeof stream?.fail === 'function'
@ -773,7 +809,10 @@ export class TextHarnessBridge {
}
async #processInteractionReply(message, messageId, senderId, key, expected) {
if (this.#signal?.aborted) return;
if (this.#signal?.aborted) {
message.statusReaction?.clear();
return;
}
const current = this.#pendingInteractions.get(key);
const claimed = expected.claimedReplyMessageId === messageId;
if (!current || current !== expected || current.submitting) {
@ -889,7 +928,10 @@ export class TextHarnessBridge {
} catch (error) {
if (error?.code === 'interaction-not-pending') {
this.#clearPendingInteraction(key, pending.interactionId);
if (this.#signal?.aborted) return;
if (this.#signal?.aborted) {
message.statusReaction?.clear();
return;
}
try {
await this.#bot.sendText(target, t(INTERACTION_RESOLVED_TEXT));
} catch (sendError) {
@ -900,7 +942,11 @@ export class TextHarnessBridge {
}
return;
}
if (this.#signal?.aborted || this.#pendingInteractions.get(key) !== pending) return;
if (this.#signal?.aborted) {
message.statusReaction?.clear();
return;
}
if (this.#pendingInteractions.get(key) !== pending) return;
pending.submitting = false;
pending.answers.pop();
pending.index -= 1;

View file

@ -20,6 +20,7 @@ oauth_config:
- files:read
- files:write
- im:history
- reactions:write
settings:
event_subscriptions:
bot_events:

View file

@ -276,6 +276,34 @@ export class SlackApi {
});
}
addReaction({ channelId, messageTs, emojiName, signal, timeoutMs }) {
return this.#request('reactions.add', {
tokenKind: 'bot',
signal,
timeoutMs,
retry: false,
body: {
channel: slackId(channelId, 'channel id'),
timestamp: requiredString(messageTs, 'message timestamp'),
name: requiredString(emojiName, 'reaction name'),
},
});
}
removeReaction({ channelId, messageTs, emojiName, signal, timeoutMs }) {
return this.#request('reactions.remove', {
tokenKind: 'bot',
signal,
timeoutMs,
retry: false,
body: {
channel: slackId(channelId, 'channel id'),
timestamp: requiredString(messageTs, 'message timestamp'),
name: requiredString(emojiName, 'reaction name'),
},
});
}
startStream({ channelId, threadTs, recipientTeamId, recipientUserId, markdownText, signal }) {
return this.#request('chat.startStream', {
tokenKind: 'bot',

View file

@ -4,6 +4,11 @@ export const SLACK_DESCRIPTOR = Object.freeze({
key: 'slack',
label: 'Slack',
connectionLabel: ' Socket Mode 长连接',
reactions: Object.freeze({
processing: 'eyes',
success: 'white_check_mark',
error: 'x',
}),
});
export class SlackHarnessBridge extends TextHarnessBridge {

View file

@ -127,6 +127,10 @@ export function normalizeSlackEvent(payload, botUserId, {
? event.files.map((file) => slackFileSource(file, loadFileStream, loadFileInfo)).filter(Boolean)
: [],
addressed: direct || mentioned,
reactionTarget: {
channelId: String(event.channel),
messageTs: String(event.ts),
},
replyTarget: {
channelId: String(event.channel),
threadTs,
@ -284,6 +288,26 @@ export class SlackBotClient {
return { providerMessageIds };
}
async addReaction(target, emoji, { signal } = {}) {
const reactionKey = String(emoji ?? '').trim();
await this.#api.addReaction({
channelId: target.channelId,
messageTs: target.messageTs,
emojiName: reactionKey,
signal: signal ?? this.#signal,
});
return reactionKey;
}
removeReaction(target, reactionKey, { signal } = {}) {
return this.#api.removeReaction({
channelId: target.channelId,
messageTs: target.messageTs,
emojiName: reactionKey,
signal: signal ?? this.#signal,
});
}
openStream(target) {
return createSlackMessageStream({
api: this.#api,

View file

@ -203,6 +203,15 @@ export class TelegramApi {
}, { signal });
}
async setMessageReaction({ chatId, messageId, emoji, signal, timeoutMs }) {
const normalizedEmoji = cleanString(emoji);
return this.#call('setMessageReaction', {
chat_id: chatId,
message_id: messageId,
reaction: normalizedEmoji ? [{ type: 'emoji', emoji: normalizedEmoji }] : [],
}, { signal, timeoutMs });
}
async sendRichMessage({
chatId,
richMessage,

View file

@ -4,6 +4,7 @@ export const TELEGRAM_DESCRIPTOR = Object.freeze({
key: 'telegram',
label: 'Telegram',
connectionLabel: ' Bot API 长轮询',
reactions: Object.freeze({ processing: '👀', success: '👍', error: '👎' }),
});
export class TelegramHarnessBridge extends TextHarnessBridge {

View file

@ -168,6 +168,7 @@ export function normalizeTelegramUpdate(update, {
images: image ? [image] : [],
files: file ? [file] : [],
addressed,
reactionTarget: { chatId, messageId },
replyTarget: {
chatId,
chatType: message.chat.type,
@ -303,6 +304,25 @@ export class TelegramBotClient {
return { providerMessageIds };
}
async addReaction(target, emoji, { signal } = {}) {
const reactionKey = String(emoji ?? '').trim();
await this.#api.setMessageReaction({
chatId: target.chatId,
messageId: target.messageId,
emoji: reactionKey,
signal: signal ?? this.#signal,
});
return reactionKey;
}
removeReaction(target, _reactionKey, { signal } = {}) {
return this.#api.setMessageReaction({
chatId: target.chatId,
messageId: target.messageId,
signal: signal ?? this.#signal,
});
}
sendTyping(target) {
return this.#api.sendChatAction({
chatId: target.chatId,

View file

@ -6,6 +6,7 @@ export const WHATSAPP_DESCRIPTOR = Object.freeze({
label: 'WhatsApp',
// Translated lazily: t() must run after setImHostLanguage, not at import time.
get connectionLabel() { return t(' Web 关联设备'); },
reactions: Object.freeze({ processing: '👀', success: '✅', error: '❌' }),
});
export class WhatsappHarnessBridge extends TextHarnessBridge {

View file

@ -246,6 +246,7 @@ export function normalizeWhatsappMessage(message, accountJid, {
addressed: !group || fromMe || mentioned || replyToSelf,
selfChat,
replyTarget: { jid: remoteJid, quoted: message, selfChat },
reactionTarget: { jid: remoteJid, key: message.key },
};
}
@ -411,6 +412,41 @@ export class WhatsappBotClient {
}, 'image');
}
async addReaction(target, emoji, { signal } = {}) {
if (typeof emoji !== 'string' || !emoji.trim()) {
throw new TypeError('A WhatsApp reaction emoji is required');
}
const reactionKey = emoji.trim();
await this.#sendReaction(target, reactionKey, signal);
return reactionKey;
}
removeReaction(target, _reactionKey, { signal } = {}) {
return this.#sendReaction(target, '', signal);
}
async #sendReaction(target, text, signal) {
if (typeof target?.jid !== 'string' || !target.jid || !target.key?.id) {
throw new TypeError('A WhatsApp reaction target is required');
}
const operationSignal = signal ?? this.#signal;
operationSignal?.throwIfAborted();
const messageId = randomBytes(10).toString('hex').toUpperCase();
if (typeof this.#outboundIds.reserve === 'function') {
this.#outboundIds.reserve(messageId);
} else {
this.#outboundIds.remember(messageId);
}
const pending = this.#socket.sendMessage(
target.jid,
{ react: { text, key: target.key } },
{ messageId },
);
const result = await waitWithSignal(pending, operationSignal);
this.#outboundIds.remember(result?.key?.id);
return result;
}
async #sendArtifact(target, file, content, presentation) {
this.#signal?.throwIfAborted();
await this.#stopTyping(target.jid);