feat: add private chat batch input commands

This commit is contained in:
xmanrui 2026-08-25 02:39:31 +08:00
parent f6cc1cae30
commit 068cf57259
28 changed files with 2665 additions and 219 deletions

View file

@ -24,6 +24,12 @@ import {
} from '../shared/preset-command.mjs';
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
import {
BatchInputManager,
batchInputBusyMessage,
batchInputGroupUnsupportedMessage,
isBatchInputCommand,
} from '../shared/batch-input.mjs';
import {
hasInboundImages,
imagePromptUserMessage,
@ -67,6 +73,9 @@ const HELP_TEXT_LINES = [
'/preset --default 跟随 Host 默认',
'/stop 停止当前任务',
'/steer 补充指令 纠偏当前任务',
'/batch 开始批量输入(仅私聊,最多 10 条文字)',
'/send 提交当前批次',
'/cancel 取消当前批次',
'/status 检查连接状态',
'/help 显示本帮助',
];
@ -344,6 +353,7 @@ export class DingtalkHarnessBridge {
#commandTasks = new Set();
#acceptedMessageIds = new Set();
#approvals;
#batchInputs = new BatchInputManager();
constructor({
api,
@ -416,12 +426,52 @@ export class DingtalkHarnessBridge {
clientSecret: this.#clientSecret,
});
const commandText = nonEmptyString(promptMessage.content) ?? '';
const addressed = String(message.conversationType) !== '2' || message?.isInAtList === true;
const direct = String(message.conversationType) !== '2';
const batchCommand = String(message?.msgtype).toLowerCase() === 'text'
&& isBatchInputCommand(commandText);
const batchStatus = this.#batchInputs.status(key);
if (batchCommand && !direct && sessionWebhook && addressed) {
return this.#finishBatchResult(
messageId,
sessionWebhook,
{ message: batchInputGroupUnsupportedMessage() },
);
}
if (direct && sessionWebhook && (batchCommand || batchStatus.phase === 'collecting')) {
const exactBatchStart = /^\/batch$/iu.test(commandText);
const result = exactBatchStart
&& batchStatus.phase === 'idle'
&& (this.#queues.has(key) || pending || this.#approvals.hasPending(key))
? { handled: true, kind: 'busy', message: batchInputBusyMessage() }
: this.#batchInputs.handle(key, commandText, {
plainText: Boolean(commandText)
&& String(message?.msgtype).toLowerCase() === 'text'
&& !hasInboundFiles(promptMessage)
&& !hasInboundImages(promptMessage),
});
if (result.handled) {
if (result.kind === 'submit') {
return this.#enqueueMessage(
{
...message,
msgtype: 'text',
text: { content: result.prompt },
},
messageId,
sender,
key,
{ batchSubmission: result },
);
}
return this.#finishBatchResult(messageId, sessionWebhook, result);
}
}
const commandRunner = hasInboundFiles(promptMessage) ? null : isControlCommand(commandText)
? runControlCommand
: (isModelCommand(commandText)
? runModelCommand
: (isPresetCommand(commandText) ? runPresetCommand : null));
const addressed = String(message.conversationType) !== '2' || message?.isInAtList === true;
if (commandRunner && sessionWebhook && addressed) {
let task;
task = this.#processFastCommand(
@ -521,6 +571,7 @@ export class DingtalkHarnessBridge {
#enqueueMessage(message, messageId, sender, key, {
releaseMessageId = true,
alreadyRecorded = false,
batchSubmission = null,
} = {}) {
let hasSafeReplyRoute = false;
try {
@ -543,6 +594,7 @@ export class DingtalkHarnessBridge {
.then(() => this.#process(message, messageId, sender, key, {
alreadyRecorded,
preparedMessage,
batchSubmission,
}))
.finally(() => {
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
@ -595,9 +647,32 @@ export class DingtalkHarnessBridge {
this.#status.lastError = null;
}
#finishBatchResult(messageId, sessionWebhook, result) {
let task;
task = Promise.resolve().then(async () => {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
increment(this.#status, 'messagesReceived');
this.#status.lastMessageAt = new Date().toISOString();
if (result.message) await this.#send(sessionWebhook, result.message);
this.#status.lastError = null;
}).catch(async (error) => {
if (this.#signal?.aborted) return;
this.#status.lastError = t('钉钉命令处理失败。');
this.#logger.error?.('[dsh-dingtalk] failed to process a batch input message', safeErrorDiagnostic(error));
await this.#send(sessionWebhook, t(CARD_ERROR_TEXT)).catch(() => undefined);
}).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
async #process(message, messageId, sender, key, {
alreadyRecorded = false,
preparedMessage,
batchSubmission = null,
} = {}) {
this.#signal?.throwIfAborted();
if (!alreadyRecorded) {
@ -633,6 +708,7 @@ export class DingtalkHarnessBridge {
const isPlainText = String(message?.msgtype).toLowerCase() === 'text';
let cardStream = null;
let cardStarted = false;
let batchSettled = batchSubmission === null;
try {
if (!text && !hasImages && !hasFiles) {
await this.#send(sessionWebhook, t('目前支持文字、图片和文件消息。'));
@ -717,6 +793,10 @@ export class DingtalkHarnessBridge {
files: promptMessage.files,
},
});
if (batchSubmission) {
this.#batchInputs.complete(key, batchSubmission.token);
batchSettled = true;
}
const answerText = typeof answer === 'string' && answer.trim()
? answer
: artifacts.length > 0 ? t('结果文件已生成。') : answer;
@ -753,6 +833,15 @@ export class DingtalkHarnessBridge {
this.#status.lastError = null;
return delivery.receipt;
} catch (error) {
let batchFailureMessage = null;
if (!batchSettled && batchSubmission) {
if (error?.code === 'turn-stopped') {
this.#batchInputs.complete(key, batchSubmission.token);
} else {
batchFailureMessage = this.#batchInputs.fail(key, batchSubmission.token).message ?? null;
}
batchSettled = true;
}
if (error?.code === 'turn-stopped') {
if (cardStarted) await cardStream.finish(t('已停止。')).catch(() => undefined);
return;
@ -767,8 +856,11 @@ export class DingtalkHarnessBridge {
const errorText = inboundFileUserMessage(error)
?? dingtalkImageErrorUserMessage(error)
?? t(CARD_ERROR_TEXT);
const streamed = cardStarted && await cardStream.finish(errorText);
if (!streamed) await this.#send(sessionWebhook, errorText);
const visibleError = batchFailureMessage
? `${errorText}\n\n${batchFailureMessage}`
: errorText;
const streamed = cardStarted && await cardStream.finish(visibleError);
if (!streamed) await this.#send(sessionWebhook, visibleError);
} catch {
this.#logger.error?.('[dsh-dingtalk] failed to send the safe error reply');
}

View file

@ -217,6 +217,9 @@ export function normalizeDiscordMessage(message, botId, { fetchImpl = fetch } =
kind: direct ? 'direct' : 'group',
conversationId: String(message.channel_id),
content: stripBotMention(message.content ?? '', botId),
plainText: (!Array.isArray(message.attachments) || message.attachments.length === 0)
&& (!Array.isArray(message.sticker_items) || message.sticker_items.length === 0)
&& !message.poll,
images: Array.isArray(message.attachments)
? message.attachments.map((attachment) => discordImageSource(attachment, fetchImpl)).filter(Boolean)
: [],

View file

@ -22,6 +22,12 @@ import {
validHarnessQuestion,
} from '../shared/harness-question.mjs';
import { HarnessApprovalQueue } from '../shared/harness-approval.mjs';
import {
BatchInputManager,
batchInputBusyMessage,
batchInputGroupUnsupportedMessage,
isBatchInputCommand,
} from '../shared/batch-input.mjs';
import { runCompactCommand } from '../shared/compact-command.mjs';
import {
isControlCommand,
@ -366,6 +372,7 @@ export class FeishuHarnessBridge {
#harness;
#state;
#queues = new Map();
#batchInputs = new BatchInputManager();
#pendingInteractions = new Map();
#interactionKeys = new Map();
#resolvedQuestionReplies = new Map();
@ -542,6 +549,58 @@ export class FeishuHarnessBridge {
const processingReaction = this.#addReaction(messageId, 'OnIt');
const commandMessage = extractInboundMessage(event, this.#client);
const commandText = nonEmptyString(commandMessage.content) ?? '';
const batchText = event.message.message_type === 'text'
? nonEmptyString(extractText(event)) ?? ''
: '';
const batchCommand = event.message.message_type === 'text'
&& isBatchInputCommand(batchText);
const pending = this.#pendingInteractions.get(key);
const batchStatus = this.#batchInputs.status(key);
if (batchCommand && event.message.chat_type !== 'p2p') {
return this.#finishBatchResult(
event,
messageId,
processingReaction,
{ message: batchInputGroupUnsupportedMessage() },
);
}
if (event.message.chat_type === 'p2p'
&& (batchCommand || batchStatus.phase === 'collecting')) {
const exactBatchStart = /^\/batch$/iu.test(batchText);
const result = exactBatchStart
&& batchStatus.phase === 'idle'
&& (this.#queues.has(key) || pending || this.#approvals.hasPending(key))
? { handled: true, kind: 'busy', message: batchInputBusyMessage() }
: this.#batchInputs.handle(key, batchText, {
plainText: event.message.message_type === 'text' && Boolean(batchText),
});
if (result.handled) {
if (result.kind === 'submit') {
const submissionEvent = {
...event,
batchSubmission: { token: result.token },
message: {
...event.message,
message_type: 'text',
content: JSON.stringify({ text: result.prompt }),
mentions: [],
},
};
return this.#enqueueMessage(
submissionEvent,
messageId,
key,
processingReaction,
);
}
return this.#finishBatchResult(
event,
messageId,
processingReaction,
result,
);
}
}
// Card commands (/m, /help, /status, etc.) bypass the queue so they
// respond immediately even when a harness task is still streaming.
if (CARD_COMMAND.test(commandText)) {
@ -599,7 +658,6 @@ export class FeishuHarnessBridge {
.finally(() => this.#acceptedMessageIds.delete(messageId));
return current;
}
const pending = this.#pendingInteractions.get(key);
const approvalReply = this.#approvals.claimReply({
key,
actor: senderOpenId(event),
@ -691,6 +749,32 @@ export class FeishuHarnessBridge {
return this.#enqueueMessage(event, messageId, key, processingReaction);
}
#finishBatchResult(event, messageId, processingReaction, result) {
let current;
current = Promise.resolve()
.then(async () => {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.lastMessageAt = new Date().toISOString();
this.#status.messagesReceived += 1;
if (result?.message) await this.#send(event.message.chat_id, result.message);
this.#status.lastError = null;
})
.then(() => this.#finishReaction(messageId, processingReaction, 'DONE'))
.catch((error) => this.#handleMessageFailure(
event,
messageId,
processingReaction,
error,
))
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(current);
});
this.#commandTasks.add(current);
return current;
}
#enqueueMessage(event, messageId, key, processingReaction, {
releaseMessageId = true,
alreadyRecorded = false,
@ -725,6 +809,9 @@ export class FeishuHarnessBridge {
async #handleMessageFailure(event, messageId, processingReaction, error) {
if (error?.code === 'turn-stopped') {
await this.#removeProcessingReaction(messageId, processingReaction);
if (error?.batchInputMessage) {
await this.#send(event.message.chat_id, error.batchInputMessage).catch(() => undefined);
}
return;
}
if (this.#signal?.aborted) {
@ -736,7 +823,8 @@ export class FeishuHarnessBridge {
await this.#finishReaction(messageId, processingReaction, 'ERROR');
await this.#send(
event.message.chat_id,
inboundFileUserMessage(error)
error?.batchInputMessage
?? inboundFileUserMessage(error)
?? imagePromptUserMessage(error)
?? t('处理失败,请稍后重试。如果问题持续,请在 DeepSeek Harness 的飞书插件页面检查连接状态。'),
).catch(() => undefined);
@ -939,12 +1027,38 @@ export class FeishuHarnessBridge {
}
this.#logger.info?.(`[dsh-feishu] processing ${event.message.chat_type} message ${messageId}`);
const batchSubmission = event.batchSubmission ?? null;
let batchAskCompleted = false;
try {
const receipt = await this.#answerWithStream(event, key, message);
const receipt = await this.#answerWithStream(event, key, message, {
onAskComplete: batchSubmission
? () => {
batchAskCompleted = this.#batchInputs.complete(
key,
batchSubmission.token,
).completed;
}
: undefined,
});
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
return receipt;
} catch (error) {
if (batchSubmission && !batchAskCompleted) {
if (error?.code === 'turn-stopped') {
this.#batchInputs.complete(key, batchSubmission.token);
throw error;
}
const failed = this.#batchInputs.fail(key, batchSubmission.token);
if (failed.retained) {
const batchError = new Error(error?.message ?? String(error), { cause: error });
batchError.code = error?.code;
batchError.batchInputMessage = failed.message;
throw batchError;
}
}
throw error;
} finally {
await this.#cancelPendingInteraction(key);
await this.#approvals.closeRoute(key);
@ -2937,10 +3051,16 @@ export class FeishuHarnessBridge {
};
}
async #answerWithStream(event, key, message) {
async #answerWithStream(event, key, message, { onAskComplete } = {}) {
const chatId = event.message.chat_id;
const messageId = event.message.message_id;
const text = message.content;
let askCompleted = false;
const markAskComplete = () => {
if (askCompleted) return;
askCompleted = true;
onAskComplete?.();
};
const content = hasInboundImages(message)
? await promptContentForMessage(message, { signal: this.#signal })
: undefined;
@ -2955,6 +3075,7 @@ export class FeishuHarnessBridge {
existsOptions: { signal: this.#signal },
askOptions: this.#interactionAskOptions(event, key, message.files),
});
markAskComplete();
let textReceipt;
let textSendError = null;
try {
@ -3009,6 +3130,7 @@ export class FeishuHarnessBridge {
existsOptions: { signal: this.#signal },
askOptions,
});
markAskComplete();
completedAnswer = completed.answer;
completedArtifacts = completed.artifacts ?? [];
await controller.setContent(answerTextForDelivery(completedAnswer, completedArtifacts));
@ -3067,6 +3189,7 @@ export class FeishuHarnessBridge {
existsOptions: { signal: this.#signal },
askOptions: this.#interactionAskOptions(event, key, message.files),
});
markAskComplete();
let textReceipt;
let textSendError = null;
try {

View file

@ -513,6 +513,11 @@ export function menuHelpText() {
'/reasoning [序号、等级ID或 --default] 查看或切换当前推理等级',
'/model [序号或完整模型ID] [推理等级ID] 查看或切换当前会话模型',
'',
'📦 批量输入(仅私聊)',
'/batch 开始批量输入(仅私聊,最多 10 条文字)',
'/send 提交当前批次',
'/cancel 取消当前批次',
'',
'🎮 任务控制',
'/stop 停止当前任务',
'/steer 指令 给 Agent 补充指令',
@ -564,6 +569,9 @@ const HELP_TEXT_COMMANDS = [
'`/reasoninglist` 或 `/reasonings` — 按序号列出当前模型可用推理等级',
'`/reasoning [序号、等级ID或 --default]` — 查看或切换当前推理等级',
'`/model [序号或完整模型ID] [推理等级ID]` — 查看或切换当前会话模型',
'`/batch` — 开启批量输入(仅私聊,最多 10 条文字)',
'`/send` — 提交当前批次',
'`/cancel` — 取消当前批次',
'`/repair` — 修复卡片按钮',
].join('\n');

View file

@ -20,6 +20,12 @@ import {
runPresetCommand,
} from '../shared/preset-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
import {
BatchInputManager,
batchInputBusyMessage,
batchInputGroupUnsupportedMessage,
isBatchInputCommand,
} from '../shared/batch-input.mjs';
import {
fetchImageBuffer,
hasInboundImages,
@ -80,6 +86,9 @@ function helpText() {
t('/preset --default 跟随 Host 默认'),
t('/stop 停止当前任务'),
t('/steer 补充指令 纠偏当前任务'),
t('/batch 开始批量输入(仅私聊,最多 10 条文字)'),
t('/send 提交当前批次'),
t('/cancel 取消当前批次'),
t('/status 检查连接状态'),
t('/help 显示本帮助'),
].join('\n');
@ -352,6 +361,7 @@ export class QqHarnessBridge {
#approvalTasks = new Set();
#commandTasks = new Set();
#approvals;
#batchInputs = new BatchInputManager();
constructor({
bot,
@ -407,14 +417,47 @@ export class QqHarnessBridge {
}
const pending = this.#pendingInteractions.get(key);
const commandText = safeText(message);
const allowed = this.#ownerUserOpenid === '*' || sender === this.#ownerUserOpenid;
const addressed = message.kind !== 'group'
|| message.rawEventType === 'GROUP_AT_MESSAGE_CREATE';
const batchCommand = isBatchInputCommand(commandText);
const batchStatus = this.#batchInputs.status(key);
if (batchCommand && allowed && addressed && message.kind === 'group') {
return this.#finishBatchResult(
message,
messageId,
{ message: batchInputGroupUnsupportedMessage() },
);
}
if (allowed && message.kind === 'c2c'
&& (batchCommand || batchStatus.phase === 'collecting')) {
const exactBatchStart = /^\/batch$/iu.test(commandText);
const result = exactBatchStart
&& batchStatus.phase === 'idle'
&& (this.#queues.has(key) || pending || this.#approvals.hasPending(key))
? { handled: true, kind: 'busy', message: batchInputBusyMessage() }
: this.#batchInputs.handle(key, commandText, {
plainText: Boolean(commandText)
&& !hasQqImageAttachments(message)
&& !hasQqFileAttachments(message),
});
if (result.handled) {
if (result.kind === 'submit') {
return this.#enqueueMessage(
{ ...message, content: result.prompt, attachments: [] },
messageId,
key,
{ batchSubmission: result },
);
}
return this.#finishBatchResult(message, messageId, result);
}
}
const commandRunner = hasQqFileAttachments(message) ? null : isControlCommand(commandText)
? runControlCommand
: (isModelCommand(commandText)
? runModelCommand
: (isPresetCommand(commandText) ? runPresetCommand : null));
const allowed = this.#ownerUserOpenid === '*' || sender === this.#ownerUserOpenid;
const addressed = message.kind !== 'group'
|| message.rawEventType === 'GROUP_AT_MESSAGE_CREATE';
if (commandRunner && allowed && addressed) {
let task;
task = this.#processFastCommand(
@ -494,6 +537,7 @@ export class QqHarnessBridge {
#enqueueMessage(message, messageId, key, {
releaseMessageId = true,
alreadyRecorded = false,
batchSubmission = null,
} = {}) {
const allowed = this.#ownerUserOpenid === '*' || message.senderId === this.#ownerUserOpenid;
const addressed = message.kind !== 'group'
@ -507,7 +551,11 @@ export class QqHarnessBridge {
const previous = this.#queues.get(key) ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#process(message, key, { alreadyRecorded, preparedMessage }))
.then(() => this.#process(message, key, {
alreadyRecorded,
preparedMessage,
batchSubmission,
}))
.finally(() => {
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
if (this.#queues.get(key) === current) this.#queues.delete(key);
@ -553,6 +601,29 @@ export class QqHarnessBridge {
this.#status.lastError = null;
}
#finishBatchResult(message, messageId, result) {
let task;
task = Promise.resolve().then(async () => {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (result.message) await this.#bot.sendText(message.replyTarget, result.message);
this.#status.lastError = null;
}).catch(async (error) => {
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-im:qq] failed to process a batch input message:', error);
await this.#bot.sendText(message.replyTarget, t('消息处理失败,请稍后重试。'))
.catch(() => undefined);
}).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
async #deliverArtifacts(target, replyTo, artifacts = [], baseReceipt = null) {
if (artifacts.length === 0) {
return { receipt: baseReceipt, failureNoticeVisible: false };
@ -589,7 +660,11 @@ export class QqHarnessBridge {
};
}
async #process(message, key, { alreadyRecorded = false, preparedMessage } = {}) {
async #process(message, key, {
alreadyRecorded = false,
preparedMessage,
batchSubmission = null,
} = {}) {
if (this.#signal?.aborted) return;
const messageId = nonEmptyString(message?.messageId);
const sender = nonEmptyString(message?.senderId);
@ -614,6 +689,7 @@ export class QqHarnessBridge {
const hasImages = hasInboundImages(promptMessage);
const hasFiles = hasInboundFiles(promptMessage);
let stream = null;
let batchSettled = batchSubmission === null;
try {
if (!text && !hasImages && !hasFiles) {
await this.#bot.sendText(target, t('目前支持文字、图片和文件消息。'));
@ -710,6 +786,10 @@ export class QqHarnessBridge {
files: promptMessage.files,
},
}));
if (batchSubmission) {
this.#batchInputs.complete(key, batchSubmission.token);
batchSettled = true;
}
} finally {
await Promise.allSettled([
this.#cancelPendingInteraction(key),
@ -771,6 +851,15 @@ export class QqHarnessBridge {
this.#status.lastError = null;
return delivery.receipt;
} catch (error) {
let batchFailureMessage = null;
if (!batchSettled && batchSubmission) {
if (error?.code === 'turn-stopped') {
this.#batchInputs.complete(key, batchSubmission.token);
} else {
batchFailureMessage = this.#batchInputs.fail(key, batchSubmission.token).message ?? null;
}
batchSettled = true;
}
if (error?.code === 'turn-stopped') {
try {
stream?.cancel?.();
@ -794,11 +883,12 @@ export class QqHarnessBridge {
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-im:qq] failed to process an inbound message:', error);
try {
const errorMessage = inboundFileUserMessage(error)
?? imagePromptUserMessage(error)
?? t('消息处理失败,请稍后重试。');
await this.#bot.sendText(
target,
inboundFileUserMessage(error)
?? imagePromptUserMessage(error)
?? t('消息处理失败,请稍后重试。'),
batchFailureMessage ? `${errorMessage}\n\n${batchFailureMessage}` : errorMessage,
);
await this.#state.markSeen(messageId);
} catch (sendError) {

View file

@ -0,0 +1,209 @@
import { t } from './i18n.mjs';
export const BATCH_INPUT_LIMIT = 10;
const BATCH_COMMAND = /^\/(batch|send|cancel)(?=$|\s)/iu;
const EXACT_BATCH_COMMAND = /^\/(batch|send|cancel)$/iu;
function commandName(text) {
if (typeof text !== 'string') return null;
return BATCH_COMMAND.exec(text.trim())?.[1]?.toLowerCase() ?? null;
}
function result(kind, message, extra = {}) {
return { handled: true, kind, ...(message ? { message } : {}), ...extra };
}
function progressMessage(count) {
if (count === BATCH_INPUT_LIMIT) {
return t(`当前已处于批量输入模式,已收集 {count}/{limit} 条。
请发送 /send 提交或 /cancel 取消。`, { count, limit: BATCH_INPUT_LIMIT });
}
return t(`当前已处于批量输入模式,已收集 {count}/{limit} 条。
完成后发送 /send,取消请发送 /cancel。`, { count, limit: BATCH_INPUT_LIMIT });
}
function submissionPrompt(messages) {
const sections = messages.map((message, index) => (
`${t('[消息 {index}]', { index: index + 1 })}\n${message}`
));
return [
t('以下是用户通过批量输入模式发送的多条内容,请按顺序作为同一次输入统一处理。'),
...sections,
].join('\n\n');
}
export function isBatchInputCommand(text) {
return commandName(text) !== null;
}
export function batchInputGroupUnsupportedMessage() {
return t('批量输入模式仅支持私聊,请在与机器人的私聊中使用。');
}
export function batchInputBusyMessage() {
return t(`当前聊天有正在运行的任务、待回答问题或待审批请求。
请先完成当前交互或发送 /stop,再使用 /batch。`);
}
export class BatchInputManager {
#batches = new Map();
status(key) {
const batch = this.#batches.get(key);
if (!batch) {
return Object.freeze({ phase: 'idle', count: 0, limit: BATCH_INPUT_LIMIT, full: false });
}
return Object.freeze({
phase: batch.phase,
count: batch.messages.length,
limit: BATCH_INPUT_LIMIT,
full: batch.messages.length === BATCH_INPUT_LIMIT,
});
}
handle(key, text, { plainText = true } = {}) {
const batch = this.#batches.get(key);
const name = commandName(text);
const exact = typeof text === 'string' ? EXACT_BATCH_COMMAND.exec(text.trim()) : null;
if (name && !exact) {
return result('invalid-command', t('用法:/{command}(不带参数)', { command: name }));
}
if (!batch) {
if (!name) return { handled: false };
if (!plainText) {
return result('unsupported-content', t('批量输入命令仅支持纯文字,请移除图片或文件后重试。'));
}
if (name === 'send') {
return result('no-batch', t('当前没有待提交的批量内容,请先发送 /batch。'));
}
if (name === 'cancel') {
return result('no-batch', t('当前没有正在进行的批量输入。'));
}
this.#batches.set(key, { phase: 'collecting', messages: [], token: null });
return result('started', t(`已进入批量输入模式,最多可发送 {limit} 条文字。
完成后发送 /send,取消请发送 /cancel。`, { limit: BATCH_INPUT_LIMIT }), {
count: 0,
limit: BATCH_INPUT_LIMIT,
});
}
if (!plainText && (batch.phase === 'collecting' || name)) {
return result('unsupported-content', t(`批量输入模式目前仅支持文字,这条消息未收录。
请继续发送文字,或使用 /send、/cancel。`), {
count: batch.messages.length,
limit: BATCH_INPUT_LIMIT,
});
}
if (batch.phase === 'submitting') {
if (name === 'send') {
return result('submitting', t('当前批次正在提交,请勿重复发送 /send。'));
}
if (name === 'cancel') {
return result('submitting', t(`批量内容已经提交,无法取消。
如需停止当前任务,请发送 /stop。`));
}
if (name === 'batch') {
return result('submitting', t('当前批次正在提交,请等待处理完成后再开启新批次。'));
}
return { handled: false };
}
if (name === 'batch') {
return result('status', progressMessage(batch.messages.length), {
count: batch.messages.length,
limit: BATCH_INPUT_LIMIT,
});
}
if (name === 'cancel') {
const count = batch.messages.length;
this.#batches.delete(key);
return result('cancelled', count === 0
? t('已取消批量输入。')
: t('已取消批量输入,共丢弃 {count} 条消息。', { count }), { count });
}
if (name === 'send') {
if (batch.messages.length === 0) {
return result('empty', t('当前批次还没有内容,请先发送文字,或使用 /cancel 取消。'), {
count: 0,
});
}
const messages = Object.freeze([...batch.messages]);
const token = Object.freeze({});
batch.phase = 'submitting';
batch.token = token;
return result('submit', null, {
token,
messages,
prompt: submissionPrompt(messages),
count: messages.length,
});
}
if (typeof text !== 'string') {
return result('unsupported-content', t(`批量输入模式目前仅支持文字,这条消息未收录。
请继续发送文字,或使用 /send、/cancel。`), {
count: batch.messages.length,
limit: BATCH_INPUT_LIMIT,
});
}
if (text.trim().startsWith('/')) {
return result('blocked-command', t('当前正在批量输入,请先发送 /send 提交或 /cancel 取消。'), {
count: batch.messages.length,
limit: BATCH_INPUT_LIMIT,
});
}
if (batch.messages.length === BATCH_INPUT_LIMIT) {
return result('full', t(`当前批次已满,这条消息未收录。
请先发送 /send 提交或 /cancel 取消,然后重新发送这条消息。`), {
count: BATCH_INPUT_LIMIT,
limit: BATCH_INPUT_LIMIT,
});
}
batch.messages.push(text);
const count = batch.messages.length;
return result('collected', count === BATCH_INPUT_LIMIT
? t('已收集 {count}/{limit} 条,当前批次已满,请发送 /send 提交或 /cancel 取消。', {
count,
limit: BATCH_INPUT_LIMIT,
})
: null, {
count,
limit: BATCH_INPUT_LIMIT,
});
}
complete(key, token) {
const batch = this.#batches.get(key);
if (!batch || batch.phase !== 'submitting' || batch.token !== token) {
return Object.freeze({ completed: false });
}
const count = batch.messages.length;
this.#batches.delete(key);
return Object.freeze({ completed: true, count });
}
fail(key, token) {
const batch = this.#batches.get(key);
if (!batch || batch.phase !== 'submitting' || batch.token !== token) {
return Object.freeze({ retained: false });
}
batch.phase = 'collecting';
batch.token = null;
const count = batch.messages.length;
return Object.freeze({
retained: true,
count,
message: t(`批量内容提交失败,已保留 {count} 条消息。
请再次发送 /send 重试或 /cancel 取消。`, { count }),
});
}
}

View file

@ -234,14 +234,15 @@ export default {
'/watch ID 关注会话(完成后推送)': '/watch ID Watch a session (push on completion)',
'/watchlist 关注列表': '/watchlist List watched sessions',
'/unwatch ID 取消关注': '/unwatch ID Stop watching a session',
'📦 批量输入(仅私聊)': '📦 Batch input (direct messages only)',
'🤖 预设 / 模型': '🤖 Presets / models',
'/models 列出模型': '/models List models',
'🎮 任务控制': '🎮 Task controls',
'/steer 指令 给 Agent 补充指令': '/steer INSTRUCTION Steer the Agent',
'**📋 卡片功能**\n\n1. 会话下拉 — 切换当前绑定会话\n2. 工作区下拉 — 切换工作区\n3. 🤖 预设下拉 — 切换 Agent 预设\n4. 🧠 模型下拉 — 切换模型\n5. 🆕 新会话 — 开启全新会话\n6. 📋 会话/关注 — 查看/绑定会话,管理关注\n7. ⏹ 停止 — 停止当前任务\n8. 📐 压缩 — 压缩当前会话上下文\n9. 补充指令 — 给 Agent 发送指令\n10. 🗄 归档切换 — 显示/隐藏归档会话\n11. 📊 状态 — 查看系统连接状态\n12. 📖 帮助 — 查看本帮助':
'**📋 Card features**\n\n1. Session dropdown — switch the bound session\n2. Workspace dropdown — switch workspace\n3. 🤖 Preset dropdown — switch Agent Preset\n4. 🧠 Model dropdown — switch model\n5. 🆕 New session — start fresh\n6. 📋 Sessions/watches — view or bind sessions and manage watches\n7. ⏹ Stop — stop the current task\n8. 📐 Compact — compact the current session context\n9. Steer task — send an instruction to the Agent\n10. 🗄 Archived toggle — show or hide archived sessions\n11. 📊 Status — view connection status\n12. 📖 Help — view this help',
'**⌨️ 文本命令**\n\n`/m` — 打开菜单卡片\n`/new` — 开启全新会话\n`/session ID` — 绑定已有会话\n`/sessionlist [工作区]` — 列出会话\n`/workspace 路径` — 切换工作区\n`/workspacelist` — 列出工作区\n`/status` — 查看连接状态\n`/compact` — 压缩上下文\n`/stop` — 停止当前任务\n`/steer 指令` — 补充指令\n`/watch ID` — 关注会话\n`/watchlist` — 关注列表\n`/unwatch ID` — 取消关注\n`/archived on/off` — 归档显隐\n`/presetlist` — 列出预设\n`/preset [序号/ID]` — 切换预设\n`/preset --default` — 跟随默认\n`/models` — 列出模型\n`/reasoninglist` 或 `/reasonings` — 按序号列出当前模型可用推理等级\n`/reasoning [序号、等级ID或 --default]` — 查看或切换当前推理等级\n`/model [序号或完整模型ID] [推理等级ID]` — 查看或切换当前会话模型\n`/repair` — 修复卡片按钮':
'**⌨️ Text commands**\n\n`/m` — open the menu card\n`/new` — start a new session\n`/session ID` — bind an existing session\n`/sessionlist [workspace]` — list sessions\n`/workspace PATH` — switch workspace\n`/workspacelist` — list workspaces\n`/status` — view connection status\n`/compact` — compact context\n`/stop` — stop the current task\n`/steer INSTRUCTION` — steer the task\n`/watch ID` — watch a session\n`/watchlist` — list watched sessions\n`/unwatch ID` — stop watching\n`/archived on/off` — show or hide archived sessions\n`/presetlist` — list presets\n`/preset [index/ID]` — switch preset\n`/preset --default` — follow default\n`/models` — list models\n`/reasoninglist` or `/reasonings` — list reasoning efforts for the current model\n`/reasoning [index, effort ID, or --default]` — show or switch reasoning effort\n`/model [index or full model ID] [reasoning effort ID]` — show or switch the current Session model\n`/repair` — repair card buttons',
'**⌨️ 文本命令**\n\n`/m` — 打开菜单卡片\n`/new` — 开启全新会话\n`/session ID` — 绑定已有会话\n`/sessionlist [工作区]` — 列出会话\n`/workspace 路径` — 切换工作区\n`/workspacelist` — 列出工作区\n`/status` — 查看连接状态\n`/compact` — 压缩上下文\n`/stop` — 停止当前任务\n`/steer 指令` — 补充指令\n`/watch ID` — 关注会话\n`/watchlist` — 关注列表\n`/unwatch ID` — 取消关注\n`/archived on/off` — 归档显隐\n`/presetlist` — 列出预设\n`/preset [序号/ID]` — 切换预设\n`/preset --default` — 跟随默认\n`/models` — 列出模型\n`/reasoninglist` 或 `/reasonings` — 按序号列出当前模型可用推理等级\n`/reasoning [序号、等级ID或 --default]` — 查看或切换当前推理等级\n`/model [序号或完整模型ID] [推理等级ID]` — 查看或切换当前会话模型\n`/batch` — 开启批量输入(仅私聊,最多 10 条文字)\n`/send` — 提交当前批次\n`/cancel` — 取消当前批次\n`/repair` — 修复卡片按钮':
'**⌨️ Text commands**\n\n`/m` — open the menu card\n`/new` — start a new session\n`/session ID` — bind an existing session\n`/sessionlist [workspace]` — list sessions\n`/workspace PATH` — switch workspace\n`/workspacelist` — list workspaces\n`/status` — view connection status\n`/compact` — compact context\n`/stop` — stop the current task\n`/steer INSTRUCTION` — steer the task\n`/watch ID` — watch a session\n`/watchlist` — list watched sessions\n`/unwatch ID` — stop watching\n`/archived on/off` — show or hide archived sessions\n`/presetlist` — list presets\n`/preset [index/ID]` — switch preset\n`/preset --default` — follow default\n`/models` — list models\n`/reasoninglist` or `/reasonings` — list reasoning efforts for the current model\n`/reasoning [index, effort ID, or --default]` — show or switch reasoning effort\n`/model [index or full model ID] [reasoning effort ID]` — show or switch the current Session model\n`/batch` — start batch input (direct messages only, up to 10 text messages)\n`/send` — submit the current batch\n`/cancel` — cancel the current batch\n`/repair` — repair card buttons',
'**💡 数字兜底**\n回复数字快速操作:\n**1**工作区列表 · **2**新会话 · **3**会话/关注\n**4**状态 · **5**修复 · **6**帮助':
'**💡 Number fallback**\nReply with a number for a quick action:\n**1** Workspace list · **2** New session · **3** Sessions/watches\n**4** Status · **5** Repair · **6** Help',
'从下方下拉选择补充指令;最后一项可自定义输入。':

View file

@ -96,4 +96,53 @@ export default {
// agent-preset.mjs
'Agent Preset 无效。': 'Invalid Agent Preset.',
// batch-input.mjs
'/batch 开始批量输入(仅私聊,最多 10 条文字)':
'/batch Start batch input (direct messages only, up to 10 text messages)',
'/send 提交当前批次': '/send Submit the current batch',
'/cancel 取消当前批次': '/cancel Cancel the current batch',
'开始批量输入(仅私聊)': 'Start batch input (direct messages only)',
'提交当前批次': 'Submit the current batch',
'取消当前批次': 'Cancel the current batch',
'当前已处于批量输入模式,已收集 {count}/{limit} 条。\n请发送 /send 提交或 /cancel 取消。':
'Batch input is already active with {count}/{limit} messages collected.\nSend /send to submit or /cancel to cancel.',
'当前已处于批量输入模式,已收集 {count}/{limit} 条。\n完成后发送 /send,取消请发送 /cancel。':
'Batch input is already active with {count}/{limit} messages collected.\nSend /send when finished or /cancel to cancel.',
'[消息 {index}]': '[Message {index}]',
'以下是用户通过批量输入模式发送的多条内容,请按顺序作为同一次输入统一处理。':
'The user sent the following messages in batch input mode. Process them in order as one input.',
'批量输入模式仅支持私聊,请在与机器人的私聊中使用。':
'Batch input is available only in direct messages. Please use it in a direct chat with the bot.',
'当前聊天有正在运行的任务、待回答问题或待审批请求。\n请先完成当前交互或发送 /stop,再使用 /batch。':
'This chat has a running task, unanswered question, or pending approval.\nFinish the current interaction or send /stop before using /batch.',
'用法:/{command}(不带参数)': 'Usage: /{command} (without arguments)',
'批量输入命令仅支持纯文字,请移除图片或文件后重试。':
'Batch input commands support text only. Remove the image or file and try again.',
'当前没有待提交的批量内容,请先发送 /batch。':
'There is no batch to submit. Send /batch first.',
'当前没有正在进行的批量输入。': 'There is no active batch input.',
'已进入批量输入模式,最多可发送 {limit} 条文字。\n完成后发送 /send,取消请发送 /cancel。':
'Batch input started. You can send up to {limit} text messages.\nSend /send when finished or /cancel to cancel.',
'批量输入模式目前仅支持文字,这条消息未收录。\n请继续发送文字,或使用 /send、/cancel。':
'Batch input currently supports text only, so this message was not collected.\nContinue with text, or use /send or /cancel.',
'当前批次正在提交,请勿重复发送 /send。':
'The current batch is being submitted. Do not send /send again.',
'批量内容已经提交,无法取消。\n如需停止当前任务,请发送 /stop。':
'The batch has already been submitted and cannot be cancelled.\nSend /stop if you need to stop the current task.',
'当前批次正在提交,请等待处理完成后再开启新批次。':
'The current batch is being submitted. Wait for it to finish before starting another batch.',
'已取消批量输入。': 'Batch input cancelled.',
'已取消批量输入,共丢弃 {count} 条消息。':
'Batch input cancelled; {count} messages were discarded.',
'当前批次还没有内容,请先发送文字,或使用 /cancel 取消。':
'The current batch is empty. Send some text first or use /cancel.',
'当前正在批量输入,请先发送 /send 提交或 /cancel 取消。':
'Batch input is active. Send /send to submit or /cancel to cancel first.',
'当前批次已满,这条消息未收录。\n请先发送 /send 提交或 /cancel 取消,然后重新发送这条消息。':
'The current batch is full, so this message was not collected.\nSend /send or /cancel first, then resend this message.',
'已收集 {count}/{limit} 条,当前批次已满,请发送 /send 提交或 /cancel 取消。':
'Collected {count}/{limit} messages. The batch is full; send /send or /cancel.',
'批量内容提交失败,已保留 {count} 条消息。\n请再次发送 /send 重试或 /cancel 取消。':
'Batch submission failed; {count} messages were retained.\nSend /send to retry or /cancel to cancel.',
};

View file

@ -19,6 +19,12 @@ import {
} from './preset-command.mjs';
import { askInWorkspaceSession } from './workspace-session.mjs';
import { HarnessApprovalQueue } from './harness-approval.mjs';
import {
BatchInputManager,
batchInputBusyMessage,
batchInputGroupUnsupportedMessage,
isBatchInputCommand,
} from './batch-input.mjs';
import {
hasInboundImages,
imagePromptUserMessage,
@ -119,6 +125,7 @@ export class TextHarnessBridge {
#approvalTasks = new Set();
#commandTasks = new Set();
#approvals;
#batches = new BatchInputManager();
constructor({
descriptor,
@ -173,7 +180,54 @@ export class TextHarnessBridge {
const key = `${kind}:${conversationId}`;
const pending = this.#pendingInteractions.get(key);
const text = cleanText(normalized.content);
const commandRunner = hasInboundFiles(normalized) ? null : isControlCommand(text)
const batchCommand = isBatchInputCommand(text);
if (batchCommand && normalized.kind === 'group' && normalized.addressed === true) {
return this.#finishLocalMessage(
normalized,
messageId,
batchInputGroupUnsupportedMessage(),
);
}
if (batchCommand && normalized.kind === 'direct') {
const exactBatch = /^\/batch$/iu.test(text);
if (exactBatch && (
this.#queues.has(key)
|| Boolean(pending)
|| this.#approvals.hasPending(key)
)) {
return this.#finishLocalMessage(normalized, messageId, batchInputBusyMessage());
}
const batch = this.#batches.handle(key, text, {
plainText: Boolean(text)
&& normalized.plainText !== false
&& !hasInboundImages(normalized)
&& !hasInboundFiles(normalized),
});
if (batch.handled) {
if (batch.kind === 'submit') {
return this.#enqueueMessage({
...normalized,
content: batch.prompt,
batchSubmission: { token: batch.token },
}, messageId, senderId, key);
}
return this.#finishLocalMessage(normalized, messageId, batch.message);
}
} else if (normalized.kind === 'direct'
&& this.#batches.status(key).phase === 'collecting') {
const batch = this.#batches.handle(key, text, {
plainText: Boolean(text)
&& normalized.plainText !== false
&& !hasInboundImages(normalized)
&& !hasInboundFiles(normalized),
});
if (batch.handled) {
return this.#finishLocalMessage(normalized, messageId, batch.message);
}
}
const collectingBatch = normalized.kind === 'direct'
&& this.#batches.status(key).phase === 'collecting';
const commandRunner = collectingBatch || hasInboundFiles(normalized) ? null : isControlCommand(text)
? runControlCommand
: (isModelCommand(text)
? runModelCommand
@ -254,6 +308,30 @@ export class TextHarnessBridge {
return this.#enqueueMessage(normalized, messageId, senderId, key);
}
#finishLocalMessage(message, messageId, reply) {
let task;
task = (async () => {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (reply) await this.#bot.sendText(message.replyTarget, reply);
this.#status.lastError = null;
})().catch(async (error) => {
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to process a batch input message:`,
error,
);
}).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
#enqueueMessage(message, messageId, senderId, key, {
releaseMessageId = true,
alreadyRecorded = false,
@ -373,6 +451,7 @@ export class TextHarnessBridge {
const target = message.replyTarget;
const text = cleanText(message.content);
const batchSubmission = message.batchSubmission;
let stream = null;
let semanticStream = false;
try {
@ -411,6 +490,9 @@ export class TextHarnessBridge {
t('/preset --default 跟随 Host 默认'),
t('/stop 停止当前任务'),
t('/steer 补充指令 纠偏当前任务'),
t('/batch 开始批量输入(仅私聊,最多 10 条文字)'),
t('/send 提交当前批次'),
t('/cancel 取消当前批次'),
t('/status 检查连接状态'),
t('/help 显示本帮助'),
].join('\n'));
@ -508,6 +590,9 @@ export class TextHarnessBridge {
files: message.files,
},
});
if (batchSubmission) {
this.#batches.complete(conversationKey, batchSubmission.token);
}
const fileOnlyCompletion = !cleanText(answer) && artifacts.length > 0;
const visibleAnswer = fileOnlyCompletion
? t(FILE_ONLY_COMPLETION_TEXT)
@ -579,7 +664,14 @@ export class TextHarnessBridge {
}
return delivery.receipt;
} catch (error) {
if (error?.code === 'turn-stopped') {
const turnStopped = error?.code === 'turn-stopped';
if (batchSubmission && turnStopped) {
this.#batches.complete(conversationKey, batchSubmission.token);
}
const failedBatch = batchSubmission && !turnStopped
? this.#batches.fail(conversationKey, batchSubmission.token)
: null;
if (turnStopped) {
if (stream) {
try {
await stream.finish(t('已停止。'));
@ -607,6 +699,23 @@ export class TextHarnessBridge {
return false;
}
};
if (failedBatch?.retained) {
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to submit a batch input:`,
error,
);
if (await presentStreamFailure(failedBatch.message)) return;
stream?.cancel?.();
try {
await this.#bot.sendText(target, failedBatch.message);
} catch (sendError) {
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to send the batch retry notice:`,
sendError,
);
}
return;
}
const imageErrorMessage = imagePromptUserMessage(error);
if (imageErrorMessage) {
if (await presentStreamFailure(imageErrorMessage)) return;

View file

@ -119,6 +119,7 @@ export function normalizeSlackEvent(payload, botUserId, {
kind: direct ? 'direct' : 'group',
conversationId: direct ? String(event.channel) : `${event.channel}:${threadTs}`,
content: stripBotMention(event.text ?? '', botUserId),
plainText: !Array.isArray(event.files) || event.files.length === 0,
images: Array.isArray(event.files)
? event.files.map((file) => slackImageSource(file, loadFile)).filter(Boolean)
: [],

View file

@ -28,6 +28,9 @@ export const TELEGRAM_COMMAND_MENU = Object.freeze([
{ command: 'preset', description: '查看或设置新会话 Agent Preset' },
{ command: 'stop', description: '停止当前任务' },
{ command: 'steer', description: '纠偏当前任务' },
{ command: 'batch', description: '开始批量输入(仅私聊)' },
{ command: 'send', description: '提交当前批次' },
{ command: 'cancel', description: '取消当前批次' },
{ command: 'status', description: '检查连接状态' },
{ command: 'help', description: '显示帮助' },
]);
@ -160,6 +163,7 @@ export function normalizeTelegramUpdate(update, {
conversationId: messageThreadId === undefined
? String(chatId) : `${chatId}:${messageThreadId}`,
content: withoutBotMention(message.text ?? message.caption ?? '', username),
plainText: typeof message.text === 'string',
images: image ? [image] : [],
files: file ? [file] : [],
addressed,

View file

@ -5,6 +5,12 @@ import {
validHarnessQuestion,
} from '../shared/harness-question.mjs';
import { HarnessApprovalQueue } from '../shared/harness-approval.mjs';
import {
BatchInputManager,
batchInputBusyMessage,
batchInputGroupUnsupportedMessage,
isBatchInputCommand,
} from '../shared/batch-input.mjs';
import { runCompactCommand } from '../shared/compact-command.mjs';
import {
isControlCommand,
@ -63,6 +69,9 @@ function helpText() {
t('/preset --default 跟随 Host 默认'),
t('/stop 停止当前任务'),
t('/steer 补充指令 纠偏当前任务'),
t('/batch 开始批量输入(仅私聊,最多 10 条文字)'),
t('/send 提交当前批次'),
t('/cancel 取消当前批次'),
t('/status 检查连接状态'),
t('/help 显示本帮助'),
].join('\n');
@ -254,6 +263,10 @@ function interactionReplyText(frame) {
return bodyOf(frame).msgtype === 'text' ? messageText(frame) : '';
}
function isNativeWecomText(frame) {
return bodyOf(frame).msgtype === 'text' && Boolean(nonEmptyString(messageText(frame)));
}
function splitUtf8(text, maxBytes = MAX_REPLY_BYTES) {
const source = String(text ?? '').trim();
if (!source) return [];
@ -466,6 +479,7 @@ export class WecomHarnessBridge {
#approvalTasks = new Set();
#commandTasks = new Set();
#approvals;
#batchInputs = new BatchInputManager();
#prefetchedImageCount = 0;
constructor({
@ -523,6 +537,40 @@ export class WecomHarnessBridge {
const pending = this.#pendingInteractions.get(key);
const commandMessage = wecomInboundMessage(frame, this.#client);
const commandText = nonEmptyString(commandMessage.content) ?? '';
const batchCommand = isBatchInputCommand(commandText);
const batchStatus = this.#batchInputs.status(key);
if (batchCommand && body.chattype === 'group') {
return this.#finishBatchResult(
frame,
messageId,
chatId,
{ message: batchInputGroupUnsupportedMessage() },
);
}
if (body.chattype === 'single'
&& (batchCommand || batchStatus.phase === 'collecting')) {
const exactBatchStart = /^\/batch$/iu.test(commandText);
const result = exactBatchStart
&& batchStatus.phase === 'idle'
&& (this.#queues.has(key) || pending || this.#approvals.hasPending(key))
? { handled: true, kind: 'busy', message: batchInputBusyMessage() }
: this.#batchInputs.handle(key, commandText, {
plainText: isNativeWecomText(frame),
});
if (result.handled) {
if (result.kind === 'submit') {
return this.#enqueueMessage({
...frame,
body: {
...body,
msgtype: 'text',
text: { content: result.prompt },
},
}, messageId, key, { batchSubmission: result });
}
return this.#finishBatchResult(frame, messageId, chatId, result);
}
}
const commandRunner = hasInboundFiles(commandMessage) ? null : isControlCommand(commandText)
? runControlCommand
: (isModelCommand(commandText)
@ -616,6 +664,7 @@ export class WecomHarnessBridge {
#enqueueMessage(frame, messageId, key, {
releaseMessageId = true,
alreadyRecorded = false,
batchSubmission = null,
} = {}) {
// WeCom image URLs expire after five minutes, while a conversation turn
// may legally stay queued longer. Start the authenticated SDK download as
@ -637,7 +686,11 @@ export class WecomHarnessBridge {
const previous = this.#queues.get(key) ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#process(frame, { alreadyRecorded, preparedMessage }))
.then(() => this.#process(frame, {
alreadyRecorded,
preparedMessage,
batchSubmission,
}))
.finally(() => {
this.#prefetchedImageCount -= reservedImages;
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
@ -647,6 +700,29 @@ export class WecomHarnessBridge {
return current;
}
#finishBatchResult(frame, messageId, chatId, result) {
let task;
task = Promise.resolve().then(async () => {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (result.message) await this.#sendImmediate(frame, chatId, result.message);
this.#status.lastError = null;
}).catch(async (error) => {
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-im:wecom] failed to process a batch input message');
await this.#sendImmediate(frame, chatId, t('消息处理失败,请稍后重试。'))
.catch(() => undefined);
}).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
async waitForIdle() {
await Promise.allSettled([
...this.#queues.values(),
@ -748,7 +824,11 @@ export class WecomHarnessBridge {
};
}
async #process(frame, { alreadyRecorded = false, preparedMessage } = {}) {
async #process(frame, {
alreadyRecorded = false,
preparedMessage,
batchSubmission = null,
} = {}) {
if (this.#signal?.aborted) return;
const body = bodyOf(frame);
const messageId = typeof body.msgid === 'string' ? body.msgid : '';
@ -767,6 +847,7 @@ export class WecomHarnessBridge {
const key = conversationKey(frame);
let streamId = null;
let streamStarted = false;
let batchSettled = batchSubmission === null;
try {
if (!text && !hasImages && !hasFiles) {
await this.#sendImmediate(frame, chatId, t('目前支持文字、图片、文件和语音转写消息。'));
@ -855,6 +936,10 @@ export class WecomHarnessBridge {
files: message.files,
},
});
if (batchSubmission) {
this.#batchInputs.complete(key, batchSubmission.token);
batchSettled = true;
}
this.#signal?.throwIfAborted();
const displayAnswer = answerTextForDelivery(answer, artifacts);
@ -915,6 +1000,15 @@ export class WecomHarnessBridge {
this.#status.lastError = null;
return delivery.receipt;
} catch (error) {
let batchFailureMessage = null;
if (!batchSettled && batchSubmission) {
if (error?.code === 'turn-stopped') {
this.#batchInputs.complete(key, batchSubmission.token);
} else {
batchFailureMessage = this.#batchInputs.fail(key, batchSubmission.token).message ?? null;
}
batchSettled = true;
}
if (error?.code === 'turn-stopped') {
if (streamStarted && streamId) {
await this.#client.replyStream(frame, streamId, t('已停止。'), true)
@ -929,11 +1023,14 @@ export class WecomHarnessBridge {
const errorText = inboundFileUserMessage(error)
?? imagePromptUserMessage(error)
?? t('消息处理失败,请稍后重试。');
const visibleError = batchFailureMessage
? `${errorText}\n\n${batchFailureMessage}`
: errorText;
try {
if (streamStarted && streamId) {
await this.#client.replyStream(frame, streamId, errorText, true);
await this.#client.replyStream(frame, streamId, visibleError, true);
} else {
await this.#sendImmediate(frame, chatId, errorText);
await this.#sendImmediate(frame, chatId, visibleError);
}
await this.#state.markSeen(messageId);
} catch {

View file

@ -12,6 +12,11 @@ import {
validHarnessQuestion,
} from '../shared/harness-question.mjs';
import { HarnessApprovalQueue } from '../shared/harness-approval.mjs';
import {
BatchInputManager,
batchInputBusyMessage,
isBatchInputCommand,
} from '../shared/batch-input.mjs';
import { runCompactCommand } from '../shared/compact-command.mjs';
import {
isControlCommand,
@ -70,6 +75,9 @@ const HELP_TEXT = () => [
t('/preset --default 跟随 Host 默认'),
t('/stop 停止当前任务'),
t('/steer 补充指令 纠偏当前任务'),
t('/batch 开始批量输入(仅私聊,最多 10 条文字)'),
t('/send 提交当前批次'),
t('/cancel 取消当前批次'),
t('/status 检查连接状态'),
t('/help 显示本帮助'),
].join('\n');
@ -104,6 +112,15 @@ function hasWeixinFileItems(message) {
&& message.item_list.some((item) => item?.file_item && typeof item.file_item === 'object');
}
function isNativeWeixinText(message) {
return Array.isArray(message?.item_list)
&& message.item_list.length > 0
&& message.item_list.every((item) => (
item?.type === 1 && typeof item.text_item?.text === 'string'
))
&& Boolean(nonEmptyString(extractWeixinText(message)));
}
function canClaimInteractionReply(message, pending) {
return pending.questions[pending.index]
&& nonEmptyString(message?.from_user_id) === pending.actor
@ -176,6 +193,7 @@ export class WeixinHarnessBridge {
#approvalTasks = new Set();
#commandTasks = new Set();
#approvals;
#batchInputs = new BatchInputManager();
constructor({
api,
@ -227,6 +245,35 @@ export class WeixinHarnessBridge {
const runId = nonEmptyString(message?.run_id) ?? undefined;
const pending = this.#pendingInteractions.get(key);
const commandText = nonEmptyString(extractWeixinText(message)) ?? '';
const batchCommand = isBatchInputCommand(commandText);
const batchStatus = this.#batchInputs.status(key);
if (sender === this.#ownerUserId
&& (batchCommand || batchStatus.phase === 'collecting')) {
const exactBatchStart = /^\/batch$/iu.test(commandText);
const result = exactBatchStart
&& batchStatus.phase === 'idle'
&& (this.#queues.has(key) || pending || this.#approvals.hasPending(key))
? { handled: true, kind: 'busy', message: batchInputBusyMessage() }
: this.#batchInputs.handle(key, commandText, {
plainText: isNativeWeixinText(message),
});
if (result.handled) {
if (result.kind === 'submit') {
return this.#enqueueMessage({
...message,
item_list: [{ type: 1, text_item: { text: result.prompt } }],
}, messageId, key, { batchSubmission: result });
}
return this.#finishBatchResult(
message,
messageId,
sender,
contextToken,
runId,
result,
);
}
}
const commandRunner = hasWeixinFileItems(message) ? null : isControlCommand(commandText)
? runControlCommand
: (isModelCommand(commandText)
@ -314,6 +361,7 @@ export class WeixinHarnessBridge {
#enqueueMessage(message, messageId, key, {
releaseMessageId = true,
alreadyRecorded = false,
batchSubmission = null,
} = {}) {
const preparedMessage = message.from_user_id === this.#ownerUserId
? prefetchInboundFiles(
@ -324,7 +372,11 @@ export class WeixinHarnessBridge {
const previous = this.#queues.get(key) ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#process(message, key, { alreadyRecorded, preparedMessage }))
.then(() => this.#process(message, key, {
alreadyRecorded,
preparedMessage,
batchSubmission,
}))
.finally(() => {
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
if (this.#queues.get(key) === current) this.#queues.delete(key);
@ -333,6 +385,33 @@ export class WeixinHarnessBridge {
return current;
}
#finishBatchResult(message, messageId, sender, contextToken, runId, result) {
let task;
task = Promise.resolve().then(async () => {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (result.message) {
await this.#send(sender, result.message, contextToken, runId);
}
this.#status.lastError = null;
this.#status.lastMessageError = null;
}).catch(async (error) => {
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#status.lastMessageError = safeMessageError(error);
this.#logger.error?.('[dsh-weixin] failed to process a batch input message:', error);
await this.#send(sender, GENERIC_PROCESSING_ERROR(), contextToken, runId)
.catch(() => undefined);
}).finally(() => {
this.#acceptedMessageIds.delete(messageId);
this.#commandTasks.delete(task);
});
this.#commandTasks.add(task);
return task;
}
async waitForIdle() {
await Promise.allSettled([
...this.#queues.values(),
@ -380,7 +459,11 @@ export class WeixinHarnessBridge {
this.#status.lastMessageError = null;
}
async #process(message, key, { alreadyRecorded = false, preparedMessage } = {}) {
async #process(message, key, {
alreadyRecorded = false,
preparedMessage,
batchSubmission = null,
} = {}) {
this.#signal?.throwIfAborted();
const messageId = weixinMessageId(message);
const sender = nonEmptyString(message?.from_user_id);
@ -398,6 +481,7 @@ export class WeixinHarnessBridge {
const contextToken = typeof message.context_token === 'string' ? message.context_token : undefined;
const runId = typeof message.run_id === 'string' ? message.run_id : undefined;
let batchSettled = batchSubmission === null;
try {
const promptMessage = preparedMessage ?? weixinInboundMessage(message, this.#api);
const text = promptMessage.content;
@ -479,6 +563,10 @@ export class WeixinHarnessBridge {
files: promptMessage.files,
},
}));
if (batchSubmission) {
this.#batchInputs.complete(key, batchSubmission.token);
batchSettled = true;
}
} finally {
await Promise.allSettled([
this.#cancelPendingInteraction(key),
@ -515,6 +603,15 @@ export class WeixinHarnessBridge {
this.#status.lastMessageError = null;
return delivery.receipt;
} catch (error) {
let batchFailureMessage = null;
if (!batchSettled && batchSubmission) {
if (error?.code === 'turn-stopped') {
this.#batchInputs.complete(key, batchSubmission.token);
} else {
batchFailureMessage = this.#batchInputs.fail(key, batchSubmission.token).message ?? null;
}
batchSettled = true;
}
if (error?.code === 'turn-stopped') {
await this.#state.markSeen(messageId);
return;
@ -529,7 +626,7 @@ export class WeixinHarnessBridge {
try {
await this.#send(
sender,
userMessage,
batchFailureMessage ? `${userMessage}\n\n${batchFailureMessage}` : userMessage,
contextToken,
runId,
);

View file

@ -73,6 +73,23 @@ function delay(ms, signal) {
});
}
export function orderWeixinMessages(messages) {
if (!Array.isArray(messages) || messages.length < 2) return Array.isArray(messages) ? messages : [];
const orderField = ['seq', 'create_time_ms'].find((field) => messages.every((message) => (
(typeof message?.[field] === 'number' && Number.isFinite(message[field]))
|| (typeof message?.[field] === 'string'
&& message[field].trim()
&& Number.isFinite(Number(message[field])))
)));
if (!orderField) return messages;
return messages
.map((message, index) => ({ message, index, order: Number(message[orderField]) }))
.sort((left, right) => (
left.order - right.order || left.index - right.index
))
.map(({ message }) => message);
}
export function createWeixinRuntimeStatus() {
return {
startedAt: null,
@ -234,7 +251,7 @@ export class WeixinRuntime {
this.#status.lastCheckedAt = Date.now();
this.#status.lastError = null;
for (const message of response?.msgs ?? []) {
for (const message of orderWeixinMessages(response?.msgs)) {
void this.#bridge.accept(message).catch((error) => {
if (signal.aborted) return;
this.#logger.error?.(

View file

@ -239,6 +239,8 @@ export function normalizeWhatsappMessage(message, accountJid, {
kind: group ? 'group' : 'direct',
conversationId: remoteJid,
content: messageText(content),
plainText: typeof content?.conversation === 'string'
|| typeof content?.extendedTextMessage?.text === 'string',
images: image ? [image] : [],
files: file ? [file] : [],
addressed: !group || fromMe || mentioned || replyToSelf,