fix: support Harness interactions across IM channels

This commit is contained in:
xmanrui 2026-08-18 14:13:10 +08:00
parent 3e7e1dcdd7
commit cd8175d341
30 changed files with 5433 additions and 555 deletions

View file

@ -238,6 +238,7 @@ export class DiscordRuntime {
status: this.#status,
logger: this.#logger,
replyTimeoutMs: this.#replyTimeoutMs,
signal: controller.signal,
});
let timer;
try {
@ -346,7 +347,16 @@ export class DiscordRuntime {
markReady();
} else if (packet.t === 'MESSAGE_CREATE') {
const message = normalizeDiscordMessage(packet.d, this.#config.platformId);
if (message) void this.#bridge?.accept(message);
const bridge = this.#bridge;
if (message && bridge) {
void bridge.accept(message).catch((error) => {
if (generation !== this.#generation || this.#stopped) return;
this.#logger.error?.(
`[dsh-im:discord] bot ${this.#config.botId} message handling failed:`,
error,
);
});
}
}
});
addSocketListener(socket, 'close', (event = {}) => {

View file

@ -1,3 +1,11 @@
import { HarnessClient } from '../weixin/harness-client.mjs';
import { HarnessClient } from '../shared/harness-client.mjs';
export class DiscordHarnessClient extends HarnessClient {}
export class DiscordHarnessClient extends HarnessClient {
constructor(options) {
super({
...options,
rpcIdPrefix: 'discord',
logPrefix: 'dsh-discord',
});
}
}

View file

@ -1,3 +1,11 @@
import { HarnessClient } from '../weixin/harness-client.mjs';
import { HarnessClient } from '../shared/harness-client.mjs';
export class QqHarnessClient extends HarnessClient {}
export class QqHarnessClient extends HarnessClient {
constructor(options) {
super({
...options,
rpcIdPrefix: 'qq',
logPrefix: 'dsh-qq',
});
}
}

View file

@ -1,6 +1,13 @@
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import {
harnessAnswerForQuestion,
harnessQuestionText,
validHarnessQuestion,
} from '../shared/harness-question.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
const HELP_TEXT = [
'QQ 机器人已连接 DeepSeek Harness。',
'',
@ -22,6 +29,17 @@ function safeText(message) {
return typeof message?.content === 'string' ? message.content.trim() : '';
}
function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
function canClaimInteractionReply(message, pending) {
return pending.questions[pending.index]
&& nonEmptyString(message?.senderId) === pending.actor
&& (message.kind !== 'group' || message.rawEventType === 'GROUP_AT_MESSAGE_CREATE')
&& nonEmptyString(safeText(message));
}
export function createQqBridgeStatus() {
return {
messagesReceived: 0,
@ -42,7 +60,11 @@ export class QqHarnessBridge {
#status;
#logger;
#replyTimeoutMs;
#signal;
#queues = new Map();
#pendingInteractions = new Map();
#interactionKeys = new Map();
#acceptedMessageIds = new Set();
constructor({
bot,
@ -52,6 +74,7 @@ export class QqHarnessBridge {
status = createQqBridgeStatus(),
logger = console,
replyTimeoutMs = 600_000,
signal,
}) {
if (!bot || typeof bot.sendText !== 'function') throw new TypeError('QQ bot client is required');
if (!ownerUserOpenid) throw new TypeError('QQ scanner identity is required');
@ -63,6 +86,7 @@ export class QqHarnessBridge {
this.#status = status;
this.#logger = logger;
this.#replyTimeoutMs = replyTimeoutMs;
this.#signal = signal;
}
get status() {
@ -70,12 +94,52 @@ export class QqHarnessBridge {
}
accept(message) {
if (this.#signal?.aborted) return Promise.resolve();
const messageId = nonEmptyString(message?.messageId);
const sender = nonEmptyString(message?.senderId);
if (!messageId || !sender || message?.senderIsBot === true
|| !['c2c', 'group'].includes(message?.kind)
|| this.#state.hasSeen(messageId)
|| this.#acceptedMessageIds.has(messageId)) return Promise.resolve();
const key = conversationKey(message);
this.#acceptedMessageIds.add(messageId);
const pending = this.#pendingInteractions.get(key);
if (pending && sender !== pending.actor) {
return this.#enqueueMessage(message, messageId, key);
}
if (pending?.submitting || pending?.claimedReplyMessageId) {
return this.#enqueueMessage(message, messageId, key);
}
if (pending) {
if (canClaimInteractionReply(message, pending)) {
pending.claimedReplyMessageId = messageId;
}
const previous = pending.queue ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#processInteractionReply(message, messageId, key, pending))
.catch((error) => this.#handleInteractionFailure(message, messageId, error))
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
if (pending.claimedReplyMessageId === messageId) pending.claimedReplyMessageId = null;
if (pending.queue === current) pending.queue = null;
});
pending.queue = current;
return current;
}
return this.#enqueueMessage(message, messageId, key);
}
#enqueueMessage(message, messageId, key, {
releaseMessageId = true,
alreadyRecorded = false,
} = {}) {
const previous = this.#queues.get(key) ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#process(message))
.then(() => this.#process(message, key, { alreadyRecorded }))
.finally(() => {
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
if (this.#queues.get(key) === current) this.#queues.delete(key);
});
this.#queues.set(key, current);
@ -83,17 +147,25 @@ export class QqHarnessBridge {
}
async waitForIdle() {
await Promise.allSettled([...this.#queues.values()]);
await Promise.allSettled([
...this.#queues.values(),
...[...this.#pendingInteractions.values()].flatMap((pending) => (
pending.queue ? [pending.queue] : []
)),
]);
}
async #process(message) {
const messageId = typeof message?.messageId === 'string' ? message.messageId : '';
const sender = typeof message?.senderId === 'string' ? message.senderId : '';
async #process(message, key, { alreadyRecorded = false } = {}) {
if (this.#signal?.aborted) return;
const messageId = nonEmptyString(message?.messageId);
const sender = nonEmptyString(message?.senderId);
if (!messageId || !sender || message.senderIsBot === true) return;
if (!['c2c', 'group'].includes(message.kind) || this.#state.hasSeen(messageId)) return;
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (!['c2c', 'group'].includes(message.kind)) return;
if (!alreadyRecorded) {
if (this.#state.hasSeen(messageId)) return;
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
}
if (this.#ownerUserOpenid !== '*' && sender !== this.#ownerUserOpenid) {
this.#status.messagesRejected += 1;
this.#status.lastRejectedAt = new Date().toISOString();
@ -116,12 +188,11 @@ export class QqHarnessBridge {
return;
}
if (command === '/status') {
await this.#harness.ensureRunning();
await this.#harness.ensureRunning({ signal: this.#signal });
await this.#bot.sendText(target, 'QQ 机器人与 DeepSeek Harness 连接正常。');
await this.#state.markSeen(messageId);
return;
}
const key = conversationKey(message);
if (command === '/new') {
await this.#state.clearSession(key);
await this.#bot.sendText(target, '已开启新会话。请发送你的问题。');
@ -146,23 +217,38 @@ export class QqHarnessBridge {
this.#logger.warn?.('[dsh-im:qq] unable to start a QQ stream; using a text reply:', error);
}
}
const { answer } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
text,
askOptions: {
timeoutMs: this.#replyTimeoutMs,
onUpdate: stream ? async (update) => {
const progress = update.type === 'text'
? update.text
: update.type === 'tool'
? `正在使用${update.name}…`
: update.text;
if (progress) await stream.update(progress);
} : undefined,
},
});
let answer;
try {
({ answer } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
text,
createOptions: { signal: this.#signal },
existsOptions: { signal: this.#signal },
askOptions: {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
onUpdate: stream ? async (update) => {
const progress = update.type === 'text'
? update.text
: update.type === 'tool'
? `正在使用${update.name}…`
: update.text;
if (progress) await stream.update(progress);
} : undefined,
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key,
actor: sender,
target,
requiresMention: message.kind === 'group',
}),
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
},
}));
} finally {
await this.#cancelPendingInteraction(key);
}
if (stream) {
try {
await stream.update(answer);
@ -179,6 +265,7 @@ export class QqHarnessBridge {
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} catch (error) {
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-im:qq] failed to process an inbound message:', error);
try {
@ -189,4 +276,269 @@ export class QqHarnessBridge {
}
}
}
async #processInteractionReply(message, messageId, key, expected) {
this.#signal?.throwIfAborted();
const current = this.#pendingInteractions.get(key);
const claimed = expected.claimedReplyMessageId === messageId;
if (!current || current !== expected || current.submitting) {
if (claimed && (!current || current !== expected)) {
return this.#discardResolvedInteractionReply(message, messageId);
}
return this.#enqueueMessage(message, messageId, key, { releaseMessageId: false });
}
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (message.kind === 'group' && message.rawEventType !== 'GROUP_AT_MESSAGE_CREATE') return;
const text = nonEmptyString(safeText(message));
if (!text) {
await this.#bot.sendText(message.replyTarget, '请用文字回答当前问题。');
return;
}
const pending = this.#pendingInteractions.get(key);
if (!pending || pending !== expected || pending.submitting) {
if (claimed && (!pending || pending !== expected)) {
await this.#bot.sendText(message.replyTarget, INTERACTION_RESOLVED_TEXT);
return;
}
return this.#enqueueMessage(message, messageId, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
pending.target = message.replyTarget;
if (pending.needsPresentation) {
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = 'QQ 交互问题发送失败。';
this.#logger.error?.('[dsh-im:qq] failed to retry an interaction question');
pending.interaction.reconnect?.();
return;
}
const presentedPending = this.#pendingInteractions.get(key);
if (!presentedPending || presentedPending !== expected || presentedPending.submitting) {
if (claimed && (!presentedPending || presentedPending !== expected)) {
await this.#bot.sendText(message.replyTarget, INTERACTION_RESOLVED_TEXT)
.catch(() => undefined);
return;
}
return this.#enqueueMessage(message, messageId, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
}
const question = pending.questions[pending.index];
if (!question) return;
pending.answers.push(harnessAnswerForQuestion(question, text));
pending.index += 1;
if (pending.index < pending.questions.length) {
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
pending.needsPresentation = true;
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = 'QQ 交互问题发送失败。';
this.#logger.error?.('[dsh-im:qq] failed to send the next interaction question');
pending.interaction.reconnect?.();
}
return;
}
pending.submitting = true;
try {
await pending.interaction.respond({
ok: true,
value: {
sessionId: pending.sessionId,
answer: { answers: pending.answers },
},
});
this.#clearPendingInteraction(key, pending.interactionId);
this.#status.lastError = null;
} catch (error) {
if (this.#signal?.aborted) return;
if (error?.code === 'interaction-not-pending') {
this.#clearPendingInteraction(key, pending.interactionId);
await this.#bot.sendText(pending.target, INTERACTION_RESOLVED_TEXT).catch(() => undefined);
return;
}
if (this.#pendingInteractions.get(key) !== pending) return;
pending.submitting = false;
pending.answers.pop();
pending.index -= 1;
this.#status.lastError = '回答提交失败。';
this.#logger.error?.('[dsh-im:qq] failed to answer a Harness interaction');
await this.#bot.sendText(pending.target, '回答提交失败,请重新发送当前问题的答案。')
.catch(() => undefined);
}
}
async #handleInteraction(interaction, {
key,
actor,
target,
requiresMention,
}) {
// Approval remains fail-closed until #5 adds an authenticated policy.
if (interaction?.kind !== 'question') return;
const questions = interaction?.payload?.questions;
const interactionId = typeof interaction?.interactionId === 'string'
? interaction.interactionId
: interaction?.rpcId;
if (typeof interaction?.rpcId !== 'string'
|| typeof interactionId !== 'string'
|| typeof interaction.sessionId !== 'string'
|| !Array.isArray(questions)
|| questions.length === 0
|| questions.some((question) => !validHarnessQuestion(question))) {
this.#logger.warn?.('[dsh-im:qq] ignored an invalid Harness question interaction');
return;
}
if (interaction.recovered === true) {
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'QQ safely cancelled an interaction left by an earlier client.',
details: {},
},
});
await this.#bot.sendText(
target,
'检测到这个 Session 中遗留的待回答问题,已安全取消并继续处理你刚才的消息。',
).catch(() => undefined);
return;
}
const existing = this.#pendingInteractions.get(key);
if (existing?.interactionId === interactionId) {
existing.interaction = interaction;
if (existing.needsPresentation) await this.#presentInteraction(existing);
return;
}
if (this.#interactionKeys.has(interactionId)) return;
if (existing) {
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'QQ is already handling another user interaction.',
details: {},
},
});
return;
}
const pending = {
kind: 'question',
interactionId,
sessionId: interaction.sessionId,
interaction,
actor,
requiresMention,
questions,
answers: [],
index: 0,
target,
queue: null,
claimedReplyMessageId: null,
presentationPromise: null,
submitting: false,
needsPresentation: true,
};
this.#pendingInteractions.set(key, pending);
this.#interactionKeys.set(interactionId, key);
await this.#presentInteraction(pending);
}
#handleInteractionResolved(resolution) {
const interactionId = resolution?.interactionId;
if (resolution?.kind !== 'question' || typeof interactionId !== 'string') return;
const key = this.#interactionKeys.get(interactionId);
if (!key) return;
this.#clearPendingInteraction(key, interactionId);
}
#presentInteraction(pending) {
if (!pending.needsPresentation) return Promise.resolve();
if (pending.presentationPromise) return pending.presentationPromise;
const question = pending.questions[pending.index];
if (!question) return Promise.resolve();
const presentation = this.#bot.sendText(
pending.target,
harnessQuestionText(
question,
pending.index,
pending.questions.length,
{ requiresMention: pending.requiresMention },
),
).then(() => {
pending.needsPresentation = false;
}).finally(() => {
if (pending.presentationPromise === presentation) pending.presentationPromise = null;
});
pending.presentationPromise = presentation;
return presentation;
}
async #discardResolvedInteractionReply(message, messageId) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
await this.#bot.sendText(message.replyTarget, INTERACTION_RESOLVED_TEXT).catch(() => undefined);
}
#takePendingInteraction(key, interactionId) {
const pending = this.#pendingInteractions.get(key);
if (!pending
|| (interactionId !== undefined && pending.interactionId !== interactionId)) return null;
this.#pendingInteractions.delete(key);
this.#interactionKeys.delete(pending.interactionId);
return pending;
}
#clearPendingInteraction(key, interactionId) {
return this.#takePendingInteraction(key, interactionId) !== null;
}
async #cancelPendingInteraction(key) {
const pending = this.#takePendingInteraction(key);
if (!pending || pending.kind !== 'question') return;
try {
await pending.interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'The QQ interaction ended before the user answered.',
details: {},
},
}, { signal: AbortSignal.timeout(5_000) });
} catch (error) {
if (error?.code !== 'interaction-not-pending') {
this.#logger.warn?.('[dsh-im:qq] failed to cancel a pending Harness interaction');
}
}
}
async #handleInteractionFailure(message, messageId, error) {
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-im:qq] failed to process an interaction reply:', error);
if (!this.#state.hasSeen(messageId)) {
await this.#state.markSeen(messageId).catch(() => undefined);
}
await this.#bot.sendText(message.replyTarget, '消息处理失败,请稍后重试。')
.catch(() => undefined);
}
}

View file

@ -101,6 +101,8 @@ export class QqRuntime {
if (!bot || typeof bot.start !== 'function' || typeof bot.stop !== 'function') {
throw new TypeError('QQ bot factory returned an invalid client');
}
const controller = new AbortController();
this.#abortController = controller;
this.#bot = bot;
this.#bridge = new QqHarnessBridge({
bot,
@ -110,6 +112,7 @@ export class QqRuntime {
status: this.#status,
logger: this.#logger,
replyTimeoutMs: this.#replyTimeoutMs,
signal: controller.signal,
});
bot.use?.(this.#typingMiddleware({
keepAlive: true,
@ -117,8 +120,6 @@ export class QqRuntime {
|| ctx?.message?.senderId === this.#config.ownerUserOpenid,
}));
const controller = new AbortController();
this.#abortController = controller;
let readyResolve;
let readyReject;
const ready = new Promise((resolve, reject) => {
@ -141,7 +142,17 @@ export class QqRuntime {
this.#logger.warn?.(`[dsh-im:qq] bot ${this.#config.botId} connection error:`, error);
}
};
const onMessage = (_ctx, message) => this.#bridge?.accept(message);
const onMessage = (_ctx, message) => {
const task = this.#bridge?.accept(message);
if (!task) return;
void task.catch((error) => {
if (controller.signal.aborted) return;
this.#logger.error?.(
`[dsh-im:qq] bot ${this.#config.botId} message handling failed:`,
error,
);
});
};
bot.on('ready', onReady);
bot.on('resumed', onReady);
bot.on('error', onError);

View file

@ -1,10 +1,23 @@
import { runWorkspaceCommand } from './workspace-command.mjs';
import { askInWorkspaceSession } from './workspace-session.mjs';
import {
harnessAnswerForQuestion,
harnessQuestionText,
validHarnessQuestion,
} from './harness-question.mjs';
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
function cleanText(value) {
return typeof value === 'string' ? value.trim() : '';
}
function canClaimInteractionReply(message, pending, senderId) {
return pending.actor === senderId
&& (message.kind !== 'group' || message.addressed === true)
&& Boolean(cleanText(message.content));
}
export function createTextBridgeStatus() {
return {
messagesReceived: 0,
@ -25,7 +38,11 @@ export class TextHarnessBridge {
#status;
#logger;
#replyTimeoutMs;
#signal;
#queues = new Map();
#pendingInteractions = new Map();
#interactionKeys = new Map();
#acceptedMessageIds = new Set();
constructor({
descriptor,
@ -35,6 +52,7 @@ export class TextHarnessBridge {
status = createTextBridgeStatus(),
logger = console,
replyTimeoutMs = 600_000,
signal,
}) {
if (!descriptor?.key || !descriptor?.label) throw new TypeError('A channel descriptor is required');
if (!bot || typeof bot.sendText !== 'function') throw new TypeError('A bot client is required');
@ -46,6 +64,7 @@ export class TextHarnessBridge {
this.#status = status;
this.#logger = logger;
this.#replyTimeoutMs = replyTimeoutMs;
this.#signal = signal;
}
get status() {
@ -53,14 +72,69 @@ export class TextHarnessBridge {
}
accept(message) {
if (this.#signal?.aborted) return Promise.resolve();
const conversationId = cleanText(message?.conversationId);
const kind = message?.kind === 'group' ? 'group' : 'direct';
const normalized = { ...message, kind, conversationId };
const messageId = cleanText(normalized.messageId);
const senderId = cleanText(normalized.senderId);
if (!messageId || !senderId || !conversationId || normalized.senderIsBot === true
|| this.#state.hasSeen(messageId) || this.#acceptedMessageIds.has(messageId)) {
return Promise.resolve();
}
this.#acceptedMessageIds.add(messageId);
const key = `${kind}:${conversationId}`;
const pending = this.#pendingInteractions.get(key);
if (pending && pending.actor !== senderId) {
return this.#enqueueMessage(normalized, messageId, senderId, key);
}
if (pending?.submitting || pending?.claimedReplyMessageId) {
return this.#enqueueMessage(normalized, messageId, senderId, key);
}
if (pending) {
if (canClaimInteractionReply(normalized, pending, senderId)) {
pending.claimedReplyMessageId = messageId;
}
const previous = pending.queue ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#processInteractionReply(
normalized,
messageId,
senderId,
key,
pending,
))
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
if (pending.queue === current) pending.queue = null;
});
pending.queue = current;
return current;
}
return this.#enqueueMessage(normalized, messageId, senderId, key);
}
#enqueueMessage(message, messageId, senderId, key, {
releaseMessageId = true,
alreadyRecorded = false,
} = {}) {
const previous = this.#queues.get(key) ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#process({ ...message, kind, conversationId }))
.then(() => this.#process(
message,
messageId,
senderId,
key,
{ alreadyRecorded },
))
.finally(() => {
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
if (this.#queues.get(key) === current) this.#queues.delete(key);
});
this.#queues.set(key, current);
@ -68,29 +142,36 @@ export class TextHarnessBridge {
}
async waitForIdle() {
await Promise.allSettled([...this.#queues.values()]);
await Promise.allSettled([
...this.#queues.values(),
...[...this.#pendingInteractions.values()].flatMap((pending) => (
pending.queue ? [pending.queue] : []
)),
]);
}
async #process(message) {
const messageId = cleanText(message.messageId);
const senderId = cleanText(message.senderId);
if (!messageId || !senderId || !message.conversationId || message.senderIsBot === true) return;
if (this.#state.hasSeen(messageId)) return;
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (message.kind === 'group' && message.addressed !== true) {
this.#status.messagesRejected += 1;
this.#status.lastRejectedAt = new Date().toISOString();
return;
async #process(message, messageId, senderId, conversationKey, {
alreadyRecorded = false,
} = {}) {
if (!alreadyRecorded) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
}
const target = message.replyTarget;
const text = cleanText(message.content);
let stream = null;
try {
this.#signal?.throwIfAborted();
if (message.kind === 'group' && message.addressed !== true) {
this.#status.messagesRejected += 1;
this.#status.lastRejectedAt = new Date().toISOString();
return;
}
if (!text) {
await this.#bot.sendText(target, '目前仅支持文字消息。');
await this.#state.markSeen(messageId);
return;
}
const command = text.toLowerCase();
@ -107,35 +188,29 @@ export class TextHarnessBridge {
'/status 检查连接状态',
'/help 显示本帮助',
].join('\n'));
await this.#state.markSeen(messageId);
return;
}
if (command === '/status') {
await this.#harness.ensureRunning();
await this.#harness.ensureRunning({ signal: this.#signal });
await this.#bot.sendText(target, `${this.#descriptor.label}机器人与 DeepSeek Harness 连接正常。`);
await this.#state.markSeen(messageId);
return;
}
const conversationKey = `${message.kind}:${message.conversationId}`;
const workspaceCommand = await runWorkspaceCommand(text, this.#harness, conversationKey);
if (workspaceCommand) {
for (const reply of workspaceCommand.messages ?? [workspaceCommand.message]) {
await this.#bot.sendText(target, reply);
}
await this.#state.markSeen(messageId);
return;
}
if (command === '/new') {
await this.#state.clearSession(conversationKey);
await this.#bot.sendText(target, '已开启新会话。请发送你的问题。');
await this.#state.markSeen(messageId);
return;
}
await this.#bot.sendTyping?.(target).catch((error) => {
this.#logger.warn?.(`[dsh-im:${this.#descriptor.key}] typing indicator failed:`, error);
});
let stream = null;
let streamFinished = false;
if (typeof this.#bot.openStream === 'function') {
try {
@ -152,13 +227,23 @@ export class TextHarnessBridge {
state: this.#state,
key: conversationKey,
text,
createOptions: this.#signal ? { signal: this.#signal } : undefined,
existsOptions: this.#signal ? { signal: this.#signal } : undefined,
askOptions: {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
onUpdate: stream ? async (update) => {
const progress = update.type === 'text' ? update.text
: update.type === 'tool' ? `正在使用${update.name}…` : update.text;
if (progress) await stream.update(progress);
} : undefined,
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key: conversationKey,
actor: senderId,
target,
requiresMention: message.kind === 'group',
}),
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
},
});
if (stream) {
@ -174,22 +259,349 @@ export class TextHarnessBridge {
}
}
if (!streamFinished) await this.#bot.sendText(target, answer);
await this.#state.markSeen(messageId);
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} catch (error) {
stream?.cancel?.();
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.(`[dsh-im:${this.#descriptor.key}] failed to process a message:`, error);
try {
await this.#bot.sendText(target, '消息处理失败,请稍后重试。');
await this.#state.markSeen(messageId);
} catch (sendError) {
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to send the safe error reply:`,
sendError,
);
}
} finally {
await this.#cancelPendingInteraction(conversationKey);
}
}
async #processInteractionReply(message, messageId, senderId, key, expected) {
if (this.#signal?.aborted) return;
const current = this.#pendingInteractions.get(key);
const claimed = expected.claimedReplyMessageId === messageId;
if (!current || current !== expected || current.submitting) {
if (claimed && (!current || current !== expected)) {
return this.#discardResolvedInteractionReply(message, messageId);
}
return this.#enqueueMessage(message, messageId, senderId, key, {
releaseMessageId: false,
});
}
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (message.kind === 'group' && message.addressed !== true) {
this.#status.messagesRejected += 1;
this.#status.lastRejectedAt = new Date().toISOString();
return;
}
const target = message.replyTarget;
const text = cleanText(message.content);
if (!text) {
try {
await this.#bot.sendText(target, '请用文字回答当前问题。');
} catch (error) {
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to reject a non-text interaction reply:`,
error,
);
}
return;
}
const pending = this.#pendingInteractions.get(key);
if (!pending || pending !== expected || pending.submitting) {
if (claimed && (!pending || pending !== expected)) {
return this.#discardResolvedInteractionReply(message, messageId, {
alreadyRecorded: true,
});
}
return this.#enqueueMessage(message, messageId, senderId, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
pending.target = target;
if (pending.needsPresentation) {
const presentationWasInFlight = pending.presentationTask !== null;
try {
await this.#presentInteraction(pending);
} catch (error) {
this.#status.lastError = `${this.#descriptor.label}交互问题发送失败。`;
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to retry an interaction question:`,
error,
);
pending.interaction.reconnect?.();
return;
}
const presented = this.#pendingInteractions.get(key);
if (!presented || presented !== expected || presented.submitting) {
if (claimed && (!presented || presented !== expected)) {
return this.#discardResolvedInteractionReply(message, messageId, {
alreadyRecorded: true,
});
}
return this.#enqueueMessage(message, messageId, senderId, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
// A reply can arrive after the platform accepted the question message but
// before its send promise settles. In that case it is already a valid
// answer. A message which itself retried a failed presentation is not.
if (!presentationWasInFlight) return;
}
const question = pending.questions[pending.index];
if (!question) return;
pending.answers.push(harnessAnswerForQuestion(question, text));
pending.index += 1;
if (pending.index < pending.questions.length) {
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
pending.needsPresentation = true;
try {
await this.#presentInteraction(pending);
} catch (error) {
this.#status.lastError = `${this.#descriptor.label}交互问题发送失败。`;
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to send the next interaction question:`,
error,
);
pending.interaction.reconnect?.();
}
return;
}
pending.submitting = true;
try {
await pending.interaction.respond({
ok: true,
value: {
sessionId: pending.sessionId,
answer: { answers: pending.answers },
},
});
this.#clearPendingInteraction(key, pending.interactionId);
this.#status.lastError = null;
} catch (error) {
if (error?.code === 'interaction-not-pending') {
this.#clearPendingInteraction(key, pending.interactionId);
if (this.#signal?.aborted) return;
try {
await this.#bot.sendText(target, INTERACTION_RESOLVED_TEXT);
} catch (sendError) {
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to send an expired interaction notice:`,
sendError,
);
}
return;
}
if (this.#signal?.aborted || this.#pendingInteractions.get(key) !== pending) return;
pending.submitting = false;
pending.answers.pop();
pending.index -= 1;
this.#status.lastError = '回答提交失败。';
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to answer a Harness interaction:`,
error,
);
try {
await this.#bot.sendText(target, '回答提交失败,请重新发送当前问题的答案。');
} catch (sendError) {
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to send an interaction retry notice:`,
sendError,
);
}
}
}
async #handleInteraction(interaction, {
key,
actor,
target,
requiresMention,
}) {
// Approval remains deliberately unanswered until #5 supplies a policy that
// can prove both the actor and the conversation allowed to decide it.
if (interaction?.kind !== 'question') return;
const questions = interaction?.payload?.questions;
const interactionId = cleanText(interaction?.interactionId) || cleanText(interaction?.rpcId);
if (!cleanText(interaction?.rpcId)
|| !interactionId
|| !cleanText(interaction?.sessionId)
|| !Array.isArray(questions)
|| questions.length === 0
|| questions.some((question) => !validHarnessQuestion(question))) {
this.#logger.warn?.(
`[dsh-im:${this.#descriptor.key}] ignored an invalid Harness question interaction`,
);
return;
}
if (interaction.recovered === true) {
await this.#respondCancellation(
interaction,
`${this.#descriptor.label} safely cancelled an interaction left by an earlier client.`,
);
try {
await this.#bot.sendText(
target,
'检测到这个 Session 中遗留的待回答问题,已安全取消并继续处理你刚才的消息。',
);
} catch (error) {
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to send an interaction recovery notice:`,
error,
);
}
return;
}
const existing = this.#pendingInteractions.get(key);
if (existing?.interactionId === interactionId) {
existing.interaction = interaction;
if (existing.needsPresentation) await this.#presentInteraction(existing);
return;
}
if (this.#interactionKeys.has(interactionId)) return;
if (existing) {
this.#logger.warn?.(
`[dsh-im:${this.#descriptor.key}] cancelled a second pending Harness question`,
);
await this.#respondCancellation(
interaction,
`${this.#descriptor.label} is already handling another user interaction.`,
);
return;
}
const pending = {
kind: 'question',
interactionId,
sessionId: interaction.sessionId,
interaction,
actor,
requiresMention,
questions,
answers: [],
index: 0,
target,
queue: null,
claimedReplyMessageId: null,
submitting: false,
needsPresentation: true,
presentationTask: null,
};
this.#pendingInteractions.set(key, pending);
this.#interactionKeys.set(interactionId, key);
await this.#presentInteraction(pending);
}
#handleInteractionResolved(resolution) {
const interactionId = cleanText(resolution?.interactionId);
if (resolution?.kind !== 'question' || !interactionId) return;
const key = this.#interactionKeys.get(interactionId);
if (!key) return;
this.#clearPendingInteraction(key, interactionId);
}
#presentInteraction(pending) {
if (pending.presentationTask) return pending.presentationTask;
const question = pending.questions[pending.index];
if (!question) return Promise.resolve();
const task = (async () => {
await this.#bot.sendText(
pending.target,
harnessQuestionText(
question,
pending.index,
pending.questions.length,
{ requiresMention: pending.requiresMention },
),
);
pending.needsPresentation = false;
})();
pending.presentationTask = task;
task.then(
() => {
if (pending.presentationTask === task) pending.presentationTask = null;
},
() => {
if (pending.presentationTask === task) pending.presentationTask = null;
},
);
return task;
}
async #discardResolvedInteractionReply(message, messageId, {
alreadyRecorded = false,
} = {}) {
if (!alreadyRecorded) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
}
try {
await this.#bot.sendText(message.replyTarget, INTERACTION_RESOLVED_TEXT);
} catch (error) {
this.#logger.error?.(
`[dsh-im:${this.#descriptor.key}] failed to send an expired interaction notice:`,
error,
);
}
}
#takePendingInteraction(key, interactionId) {
const pending = this.#pendingInteractions.get(key);
if (!pending
|| (interactionId !== undefined && pending.interactionId !== interactionId)) return null;
this.#pendingInteractions.delete(key);
this.#interactionKeys.delete(pending.interactionId);
return pending;
}
#clearPendingInteraction(key, interactionId) {
return this.#takePendingInteraction(key, interactionId) !== null;
}
async #respondCancellation(interaction, message) {
try {
await interaction.respond({
ok: false,
error: { code: 'cancelled', message, details: {} },
}, { signal: AbortSignal.timeout(5_000) });
} catch (error) {
if (error?.code !== 'interaction-not-pending') throw error;
}
}
async #cancelPendingInteraction(key) {
const pending = this.#takePendingInteraction(key);
if (!pending || pending.kind !== 'question') return;
try {
await this.#respondCancellation(
pending.interaction,
`The ${this.#descriptor.label} interaction ended before the user answered.`,
);
} catch (error) {
this.#logger.warn?.(
`[dsh-im:${this.#descriptor.key}] failed to cancel a pending Harness interaction:`,
error,
);
}
}
}

View file

@ -1,3 +1,11 @@
import { HarnessClient } from '../weixin/harness-client.mjs';
import { HarnessClient } from '../shared/harness-client.mjs';
export class SlackHarnessClient extends HarnessClient {}
export class SlackHarnessClient extends HarnessClient {
constructor(options) {
super({
...options,
rpcIdPrefix: 'slack',
logPrefix: 'dsh-slack',
});
}
}

View file

@ -309,6 +309,7 @@ export class SlackRuntime {
status: this.#status,
logger: this.#logger,
replyTimeoutMs: this.#replyTimeoutMs,
signal: controller.signal,
});
let timer;
try {
@ -389,7 +390,16 @@ export class SlackRuntime {
if (this.#appId && packet.payload.api_app_id
&& packet.payload.api_app_id !== this.#appId) return;
const message = normalizeSlackEvent(packet.payload, this.#config.platformId.split(':')[1]);
if (message) void this.#bridge?.accept(message);
const bridge = this.#bridge;
if (message && bridge) {
void bridge.accept(message).catch((error) => {
if (generation !== this.#generation || this.#stopped) return;
this.#logger.error?.(
`[dsh-im:slack] bot ${this.#config.botId} message handling failed:`,
error,
);
});
}
});
addSocketListener(socket, 'close', (event = {}) => {

View file

@ -1,3 +1,11 @@
import { HarnessClient } from '../weixin/harness-client.mjs';
import { HarnessClient } from '../shared/harness-client.mjs';
export class TelegramHarnessClient extends HarnessClient {}
export class TelegramHarnessClient extends HarnessClient {
constructor(options) {
super({
...options,
rpcIdPrefix: 'telegram',
logPrefix: 'dsh-telegram',
});
}
}

View file

@ -34,19 +34,21 @@ export function normalizeTelegramUpdate(update, { botId, username }) {
const addressed = direct
|| String(message.reply_to_message?.from?.id ?? '') === String(botId)
|| mentionedUsername(message, username);
const messageThreadId = Number.isSafeInteger(message.message_thread_id)
? message.message_thread_id : undefined;
return {
messageId: String(update.update_id),
senderId: String(senderId),
senderIsBot: message.from?.is_bot === true,
kind: direct ? 'direct' : 'group',
conversationId: String(chatId),
conversationId: messageThreadId === undefined
? String(chatId) : `${chatId}:${messageThreadId}`,
content: withoutBotMention(message.text ?? message.caption ?? '', username),
addressed,
replyTarget: {
chatId,
replyToMessageId: messageId,
messageThreadId: Number.isSafeInteger(message.message_thread_id)
? message.message_thread_id : undefined,
messageThreadId,
},
};
}
@ -206,6 +208,7 @@ export class TelegramRuntime {
status: this.#status,
logger: this.#logger,
replyTimeoutMs: this.#replyTimeoutMs,
signal: controller.signal,
});
let cursor = this.#state.cursor();
@ -248,7 +251,15 @@ export class TelegramRuntime {
botId: this.#config.platformId,
username: this.#config.username,
});
if (message) await this.#bridge.accept(message);
if (message) {
void this.#bridge.accept(message).catch((error) => {
if (signal.aborted) return;
this.#logger.error?.(
`[dsh-im:telegram] bot ${this.#config.botId} message handling failed:`,
error,
);
});
}
cursor = update.update_id + 1;
await this.#state.setCursor(cursor);
}

View file

@ -1,3 +1,11 @@
import { HarnessClient } from '../weixin/harness-client.mjs';
import { HarnessClient } from '../shared/harness-client.mjs';
export class WecomHarnessClient extends HarnessClient {}
export class WecomHarnessClient extends HarnessClient {
constructor(options) {
super({
...options,
rpcIdPrefix: 'wecom',
logPrefix: 'dsh-wecom',
});
}
}

View file

@ -1,4 +1,9 @@
import { generateReqId } from '@wecom/aibot-node-sdk';
import {
harnessAnswerForQuestion,
harnessQuestionText,
validHarnessQuestion,
} from '../shared/harness-question.mjs';
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
@ -15,6 +20,11 @@ const HELP_TEXT = [
'/help 显示本帮助',
].join('\n');
const MAX_REPLY_BYTES = 18_000;
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
function bodyOf(frame) {
return frame?.body && typeof frame.body === 'object' ? frame.body : {};
@ -27,16 +37,23 @@ function conversationKey(frame) {
function messageText(frame) {
const body = bodyOf(frame);
if (body.msgtype === 'text') return typeof body.text?.content === 'string' ? body.text.content.trim() : '';
if (body.msgtype === 'voice') return typeof body.voice?.content === 'string' ? body.voice.content.trim() : '';
if (body.msgtype === 'mixed' && Array.isArray(body.mixed?.msg_item)) {
return body.mixed.msg_item
let text = '';
if (body.msgtype === 'text') {
text = typeof body.text?.content === 'string' ? body.text.content.trim() : '';
} else if (body.msgtype === 'voice') {
text = typeof body.voice?.content === 'string' ? body.voice.content.trim() : '';
} else if (body.msgtype === 'mixed' && Array.isArray(body.mixed?.msg_item)) {
text = body.mixed.msg_item
.filter((item) => item?.msgtype === 'text' && typeof item.text?.content === 'string')
.map((item) => item.text.content)
.join('\n')
.trim();
}
return '';
// Group callbacks retain the leading @bot mention that caused delivery.
// It is routing metadata rather than part of the user's prompt or answer.
return body.chattype === 'group'
? text.replace(/^\s*@\S+(?:\s+|$)/u, '').trim()
: text;
}
function splitUtf8(text, maxBytes = MAX_REPLY_BYTES) {
@ -66,6 +83,12 @@ function progressText(update) {
return update?.text;
}
function canClaimInteractionReply(frame, pending) {
return pending.questions[pending.index]
&& nonEmptyString(bodyOf(frame).from?.userid) === pending.actor
&& nonEmptyString(messageText(frame));
}
export function createWecomBridgeStatus() {
return {
messagesReceived: 0,
@ -86,7 +109,11 @@ export class WecomHarnessBridge {
#logger;
#replyTimeoutMs;
#generateReqId;
#signal;
#queues = new Map();
#pendingInteractions = new Map();
#interactionKeys = new Map();
#acceptedMessageIds = new Set();
constructor({
client,
@ -96,6 +123,7 @@ export class WecomHarnessBridge {
logger = console,
replyTimeoutMs = 600_000,
generateStreamId = generateReqId,
signal,
}) {
if (!client || typeof client.replyStream !== 'function' || typeof client.sendMessage !== 'function') {
throw new TypeError('Enterprise WeChat client is required');
@ -108,6 +136,7 @@ export class WecomHarnessBridge {
this.#logger = logger;
this.#replyTimeoutMs = replyTimeoutMs;
this.#generateReqId = generateStreamId;
this.#signal = signal;
}
get status() {
@ -115,12 +144,65 @@ export class WecomHarnessBridge {
}
accept(frame) {
if (this.#signal?.aborted) return Promise.resolve();
const body = bodyOf(frame);
const messageId = nonEmptyString(body.msgid);
const senderId = nonEmptyString(body.from?.userid);
const chatId = body.chattype === 'group'
? nonEmptyString(body.chatid)
: senderId;
if (!messageId || !senderId || !chatId
|| !['single', 'group'].includes(body.chattype)
|| this.#state.hasSeen(messageId)
|| this.#acceptedMessageIds.has(messageId)) return Promise.resolve();
const key = conversationKey(frame);
this.#acceptedMessageIds.add(messageId);
const pending = this.#pendingInteractions.get(key);
if (pending && pending.actor !== senderId) {
return this.#enqueueMessage(frame, messageId, key);
}
if (pending?.submitting || pending?.claimedReplyMessageId) {
return this.#enqueueMessage(frame, messageId, key);
}
if (pending) {
if (canClaimInteractionReply(frame, pending)) {
pending.claimedReplyMessageId = messageId;
}
const previous = pending.queue ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#processInteractionReply(
frame,
messageId,
senderId,
chatId,
key,
pending,
))
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
if (pending.queue === current) pending.queue = null;
});
pending.queue = current;
return current;
}
return this.#enqueueMessage(frame, messageId, key);
}
#enqueueMessage(frame, messageId, key, {
releaseMessageId = true,
alreadyRecorded = false,
} = {}) {
const previous = this.#queues.get(key) ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#process(frame))
.then(() => this.#process(frame, { alreadyRecorded }))
.finally(() => {
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
if (this.#queues.get(key) === current) this.#queues.delete(key);
});
this.#queues.set(key, current);
@ -128,16 +210,23 @@ export class WecomHarnessBridge {
}
async waitForIdle() {
await Promise.allSettled([...this.#queues.values()]);
await Promise.allSettled([
...this.#queues.values(),
...[...this.#pendingInteractions.values()].flatMap((pending) => (
pending.queue ? [pending.queue] : []
)),
]);
}
async #sendActive(chatId, text) {
for (const chunk of splitUtf8(text)) {
this.#signal?.throwIfAborted();
await this.#client.sendMessage(chatId, { msgtype: 'markdown', markdown: { content: chunk } });
}
}
async #sendImmediate(frame, chatId, text) {
this.#signal?.throwIfAborted();
const chunks = splitUtf8(text);
if (chunks.length === 0) return;
try {
@ -150,17 +239,20 @@ export class WecomHarnessBridge {
}
}
async #process(frame) {
async #process(frame, { alreadyRecorded = false } = {}) {
if (this.#signal?.aborted) return;
const body = bodyOf(frame);
const messageId = typeof body.msgid === 'string' ? body.msgid : '';
const senderId = typeof body.from?.userid === 'string' ? body.from.userid : '';
const chatId = body.chattype === 'group' ? body.chatid : senderId;
if (!messageId || !senderId || !chatId || !['single', 'group'].includes(body.chattype)) return;
if (this.#state.hasSeen(messageId)) return;
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (!alreadyRecorded) {
if (this.#state.hasSeen(messageId)) return;
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
}
const text = messageText(frame);
const key = conversationKey(frame);
let streamId = null;
let streamStarted = false;
try {
@ -176,12 +268,11 @@ export class WecomHarnessBridge {
return;
}
if (command === '/status') {
await this.#harness.ensureRunning();
await this.#harness.ensureRunning({ signal: this.#signal });
await this.#sendImmediate(frame, chatId, '企业微信机器人与 DeepSeek Harness 连接正常。');
await this.#state.markSeen(messageId);
return;
}
const key = conversationKey(frame);
if (command === '/new') {
await this.#state.clearSession(key);
await this.#sendImmediate(frame, chatId, '已开启新会话。请发送你的问题。');
@ -210,14 +301,24 @@ export class WecomHarnessBridge {
state: this.#state,
key,
text,
createOptions: { signal: this.#signal },
existsOptions: { signal: this.#signal },
askOptions: {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
onUpdate: streamStarted && typeof this.#client.replyStreamNonBlocking === 'function'
? async (update) => {
const progress = splitUtf8(progressText(update))[0];
if (progress) await this.#client.replyStreamNonBlocking(frame, streamId, progress, false);
}
: undefined,
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key,
actor: senderId,
chatId,
requiresMention: body.chattype === 'group',
}),
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
},
});
@ -240,6 +341,7 @@ export class WecomHarnessBridge {
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} catch (error) {
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-im:wecom] failed to process an inbound message');
try {
@ -252,6 +354,276 @@ export class WecomHarnessBridge {
} catch {
this.#logger.error?.('[dsh-im:wecom] failed to send the safe error reply');
}
} finally {
await this.#cancelPendingInteraction(key);
}
}
async #processInteractionReply(frame, messageId, senderId, chatId, key, expected) {
if (this.#signal?.aborted) return;
const current = this.#pendingInteractions.get(key);
const claimed = expected.claimedReplyMessageId === messageId;
if (!current || current !== expected || current.submitting) {
if (claimed && (!current || current !== expected)) {
return this.#discardResolvedInteractionReply(frame, messageId, chatId);
}
return this.#enqueueMessage(frame, messageId, key, { releaseMessageId: false });
}
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
const text = nonEmptyString(messageText(frame));
if (!text) {
await this.#sendImmediate(frame, chatId, '请用文字或语音回答当前问题。')
.catch(() => undefined);
return;
}
const pending = this.#pendingInteractions.get(key);
if (!pending || pending !== expected || pending.submitting) {
if (claimed && (!pending || pending !== expected)) {
await this.#sendImmediate(frame, chatId, INTERACTION_RESOLVED_TEXT)
.catch(() => undefined);
return;
}
return this.#enqueueMessage(frame, messageId, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
if (pending.actor !== senderId) {
return this.#enqueueMessage(frame, messageId, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
pending.chatId = chatId;
if (pending.needsPresentation) {
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = '企业微信交互问题发送失败。';
this.#logger.error?.('[dsh-im:wecom] failed to retry an interaction question');
pending.interaction.reconnect?.();
return;
}
const presentedPending = this.#pendingInteractions.get(key);
if (!presentedPending || presentedPending !== expected || presentedPending.submitting) {
if (claimed && (!presentedPending || presentedPending !== expected)) {
await this.#sendImmediate(frame, chatId, INTERACTION_RESOLVED_TEXT)
.catch(() => undefined);
return;
}
return this.#enqueueMessage(frame, messageId, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
}
const question = pending.questions[pending.index];
if (!question) return;
pending.answers.push(harnessAnswerForQuestion(question, text));
pending.index += 1;
if (pending.index < pending.questions.length) {
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
pending.needsPresentation = true;
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = '企业微信交互问题发送失败。';
this.#logger.error?.('[dsh-im:wecom] failed to send the next interaction question');
pending.interaction.reconnect?.();
}
return;
}
pending.submitting = true;
try {
await pending.interaction.respond({
ok: true,
value: {
sessionId: pending.sessionId,
answer: { answers: pending.answers },
},
});
this.#clearPendingInteraction(key, pending.interactionId);
this.#status.lastError = null;
} catch (error) {
if (this.#signal?.aborted) return;
if (error?.code === 'interaction-not-pending') {
if (this.#pendingInteractions.get(key) === pending) {
this.#clearPendingInteraction(key, pending.interactionId);
}
await this.#sendImmediate(frame, chatId, INTERACTION_RESOLVED_TEXT)
.catch(() => undefined);
return;
}
if (this.#pendingInteractions.get(key) !== pending) return;
pending.submitting = false;
pending.answers.pop();
pending.index -= 1;
this.#status.lastError = '回答提交失败。';
this.#logger.error?.('[dsh-im:wecom] failed to answer a Harness interaction');
await this.#sendImmediate(frame, chatId, '回答提交失败,请重新发送当前问题的答案。')
.catch(() => undefined);
}
}
async #handleInteraction(interaction, {
key,
actor,
chatId,
requiresMention,
}) {
// Approval remains unanswered until #5 supplies an authenticated policy.
if (interaction?.kind !== 'question') return;
const questions = interaction?.payload?.questions;
const interactionId = typeof interaction?.interactionId === 'string'
? interaction.interactionId
: interaction?.rpcId;
if (typeof interaction?.rpcId !== 'string'
|| typeof interactionId !== 'string'
|| typeof interaction.sessionId !== 'string'
|| !Array.isArray(questions)
|| questions.length === 0
|| questions.some((question) => !validHarnessQuestion(question))) {
this.#logger.warn?.('[dsh-im:wecom] ignored an invalid Harness question interaction');
return;
}
if (interaction.recovered === true) {
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'Enterprise WeChat safely cancelled an interaction left by an earlier client.',
details: {},
},
});
await this.#sendActive(
chatId,
'检测到这个 Session 中遗留的待回答问题,已安全取消并继续处理你刚才的消息。',
).catch(() => undefined);
return;
}
const existing = this.#pendingInteractions.get(key);
if (existing?.interactionId === interactionId) {
existing.interaction = interaction;
if (existing.needsPresentation) await this.#presentInteraction(existing);
return;
}
if (this.#interactionKeys.has(interactionId)) return;
if (existing) {
this.#logger.warn?.('[dsh-im:wecom] cancelled a second pending Harness question');
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'Enterprise WeChat is already handling another user interaction.',
details: {},
},
});
return;
}
const pending = {
kind: 'question',
interactionId,
sessionId: interaction.sessionId,
interaction,
actor,
requiresMention,
questions,
answers: [],
index: 0,
chatId,
queue: null,
claimedReplyMessageId: null,
submitting: false,
needsPresentation: true,
presentationPromise: null,
};
this.#pendingInteractions.set(key, pending);
this.#interactionKeys.set(pending.interactionId, key);
await this.#presentInteraction(pending);
}
#handleInteractionResolved(resolution) {
const interactionId = resolution?.interactionId;
if (resolution?.kind !== 'question' || typeof interactionId !== 'string') return;
const key = this.#interactionKeys.get(interactionId);
if (!key) return;
this.#clearPendingInteraction(key, interactionId);
}
#presentInteraction(pending) {
if (!pending.needsPresentation) return Promise.resolve();
if (pending.presentationPromise) return pending.presentationPromise;
const question = pending.questions[pending.index];
if (!question) return Promise.resolve();
const presentation = this.#sendActive(
pending.chatId,
harnessQuestionText(
question,
pending.index,
pending.questions.length,
{ requiresMention: pending.requiresMention },
),
).then(() => {
pending.needsPresentation = false;
}).finally(() => {
if (pending.presentationPromise === presentation) {
pending.presentationPromise = null;
}
});
pending.presentationPromise = presentation;
return presentation;
}
async #discardResolvedInteractionReply(frame, messageId, chatId) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
await this.#sendImmediate(frame, chatId, INTERACTION_RESOLVED_TEXT).catch(() => undefined);
}
#takePendingInteraction(key, interactionId) {
const pending = this.#pendingInteractions.get(key);
if (!pending
|| (interactionId !== undefined && pending.interactionId !== interactionId)) return null;
this.#pendingInteractions.delete(key);
this.#interactionKeys.delete(pending.interactionId);
return pending;
}
#clearPendingInteraction(key, interactionId) {
return this.#takePendingInteraction(key, interactionId) !== null;
}
async #cancelPendingInteraction(key) {
const pending = this.#takePendingInteraction(key);
if (!pending || pending.kind !== 'question') return;
try {
await pending.interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'The Enterprise WeChat interaction ended before the user answered.',
details: {},
},
}, { signal: AbortSignal.timeout(5_000) });
} catch (error) {
if (error?.code !== 'interaction-not-pending') {
this.#logger.warn?.('[dsh-im:wecom] failed to cancel a pending Harness interaction');
}
}
}
}

View file

@ -36,6 +36,7 @@ export class WecomRuntime {
#bridge = null;
#starting = null;
#startController = null;
#runtimeController = null;
constructor({
config,
@ -69,8 +70,10 @@ export class WecomRuntime {
async start() {
if (this.#status.ready && this.#client) return this.status;
if (this.#starting) return this.#starting;
this.#runtimeController?.abort(new DOMException('Enterprise WeChat runtime replaced', 'AbortError'));
const controller = new AbortController();
this.#startController = controller;
this.#runtimeController = controller;
this.#starting = this.#start(controller.signal).finally(() => {
if (this.#startController === controller) this.#startController = null;
this.#starting = null;
@ -105,6 +108,7 @@ export class WecomRuntime {
status: this.#status,
logger: this.#logger,
replyTimeoutMs: this.#replyTimeoutMs,
signal,
});
let readyResolve;
@ -186,6 +190,8 @@ export class WecomRuntime {
async stop() {
const starting = this.#starting;
this.#startController?.abort(new DOMException('Enterprise WeChat runtime stopped', 'AbortError'));
this.#runtimeController?.abort(new DOMException('Enterprise WeChat runtime stopped', 'AbortError'));
this.#runtimeController = null;
await this.#stopActive();
await starting?.catch(() => undefined);
return this.status;

View file

@ -1,329 +1,19 @@
import { spawn } from 'node:child_process';
import { randomUUID } from 'node:crypto';
import { isAbsolute } from 'node:path';
import {
HarnessClient as SharedHarnessClient,
} from '../shared/harness-client.mjs';
import { adoptRegisteredWorkspaceSession } from '../shared/harness-session-binding.mjs';
export {
HarnessInteractionError,
HarnessReplyTracker,
HarnessRpcError,
} from '../shared/harness-client.mjs';
const sleep = (ms) => new Promise((resolve) => setTimeout(resolve, ms));
function workspacePaths(value) {
if (!Array.isArray(value?.items)) return [];
return value.items.flatMap((item) => (
typeof item?.path === 'string' && isAbsolute(item.path) ? [item.path] : []
));
}
function workspaceFromList(workspacePath, workspaceList) {
if (!Array.isArray(workspaceList?.items)
|| !Array.isArray(workspaceList?.archivedSessionIds)) {
throw new Error('Harness returned an invalid response for workspace.list');
}
const workspace = workspaceList.items.find((item) => item?.path === workspacePath);
if (!workspace) return null;
if (!Array.isArray(workspace.sessionIds)
|| workspace.sessionIds.some((sessionId) => typeof sessionId !== 'string')) {
throw new Error('Harness returned invalid session IDs for workspace.list');
}
return workspace;
}
function workspaceSessions(workspace, archivedSessionIds, sessionList) {
if (!Array.isArray(sessionList?.items)) {
throw new Error('Harness returned an invalid response for session.list');
}
const archived = new Set(archivedSessionIds);
const summaries = new Map(sessionList.items.flatMap((item) => (
typeof item?.sessionId === 'string' ? [[item.sessionId, item]] : []
)));
return {
workspace: workspace.path,
sessions: workspace.sessionIds.map((sessionId) => {
const summary = summaries.get(sessionId);
const title = summary?.projections?.values?.title;
return {
sessionId,
title: typeof title === 'string' ? title : null,
archived: archived.has(sessionId),
blank: summary?.blank === true,
origin: summary?.origin === 'subagent' ? 'subagent' : null,
summaryAvailable: summary !== undefined,
};
}),
};
}
function assistantMessageText(event) {
return (event?.data?.message?.content ?? [])
.filter((part) => part.type === 'text' && typeof part.text === 'string')
.map((part) => part.text)
.join('\n')
.trim();
}
export class HarnessReplyTracker {
#promptRpcId;
#lastSeq;
#openTurn = null;
#targetTurn = null;
#stepText = new Map();
#latestText = '';
#finished = false;
#reason = null;
constructor({ promptRpcId, afterSeq = -1 }) {
this.#promptRpcId = promptRpcId;
this.#lastSeq = afterSeq;
}
get finished() {
return this.#finished;
}
get answer() {
return this.#latestText.trim();
}
get reason() {
return this.#reason;
}
consume(entries) {
let update = null;
const ordered = [...entries]
.map((entry) => entry?.event ?? entry)
.filter(Boolean)
.sort((left, right) => (left.seq ?? -1) - (right.seq ?? -1));
for (const event of ordered) {
const seq = event.seq ?? -1;
if (seq <= this.#lastSeq) continue;
this.#lastSeq = seq;
if (event.type === 'turn/start') this.#openTurn = event.data?.turn ?? null;
if (event.type === 'user/message' && event.data?.source?.rpcId === this.#promptRpcId) {
this.#targetTurn = this.#openTurn;
continue;
}
if (this.#targetTurn === null) continue;
if (event.type === 'turn/end') {
if (event.data?.turn !== this.#targetTurn) continue;
this.#finished = true;
this.#reason = event.data?.reason ?? null;
this.#openTurn = null;
continue;
}
if (event.data?.turn !== this.#targetTurn) continue;
if (event.type === 'assistant/chunk' && event.data?.chunk?.type === 'text-delta') {
const step = event.data?.step ?? 0;
const index = event.data.chunk.index ?? 0;
const key = `${step}:${index}`;
this.#stepText.set(key, (this.#stepText.get(key) ?? '') + event.data.chunk.text);
const prefix = `${step}:`;
const text = [...this.#stepText.entries()]
.filter(([partKey]) => partKey.startsWith(prefix))
.sort(([left], [right]) => Number(left.split(':')[1]) - Number(right.split(':')[1]))
.map(([, part]) => part)
.join('\n')
.trim();
if (text && text !== this.#latestText) {
this.#latestText = text;
update = { type: 'text', text };
}
continue;
}
if (event.type === 'assistant/message') {
const text = assistantMessageText(event);
if (text && text !== this.#latestText) {
this.#latestText = text;
update = { type: 'text', text };
}
continue;
}
if (event.type === 'tool/call') {
update = { type: 'tool', name: event.data?.name ?? '工具' };
} else if (event.type === 'tool/result') {
update = { type: 'status', text: '正在整理结果…' };
}
}
return update;
}
}
export class HarnessRpcError extends Error {
constructor(method, error) {
super(`${method}: ${error?.message ?? 'unknown Harness RPC error'}`);
this.name = 'HarnessRpcError';
this.method = method;
this.code = error?.code ?? 'internal';
this.details = error?.details ?? {};
}
}
export class HarnessClient {
#baseUrl;
#workspace;
#agentPreset;
#autostart;
#dshBin;
#managedProcess = null;
constructor({ baseUrl, workspace, agentPreset = 'standard', autostart = false, dshBin = 'dsh' }) {
this.#baseUrl = new URL(baseUrl);
this.#workspace = workspace;
this.#agentPreset = agentPreset;
this.#autostart = autostart;
this.#dshBin = dshBin;
}
async rpc(method, payload = {}, timeoutMs = 30_000, options = {}) {
const rpcId = options.rpcId ?? `weixin-${randomUUID()}`;
const response = await fetch(new URL(`/api/${method}`, this.#baseUrl), {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ type: 'client-request', rpcId, method, payload }),
signal: AbortSignal.timeout(timeoutMs),
export class HarnessClient extends SharedHarnessClient {
constructor(options) {
super({
...options,
rpcIdPrefix: 'weixin',
logPrefix: 'dsh-weixin',
});
if (!response.ok) throw new Error(`Harness transport ${method} failed: HTTP ${response.status}`);
const body = await response.json();
if (body?.type !== 'server-response' || body?.rpcId !== rpcId) {
throw new Error(`Harness returned an invalid response for ${method}`);
}
if (!body.result?.ok) throw new HarnessRpcError(method, body.result?.error);
return body.result.value;
}
async health() {
await this.rpc('host.describe', {}, 5_000);
return true;
}
async ensureRunning() {
try {
return await this.health();
} catch (firstError) {
if (!this.#autostart) throw firstError;
}
if (!this.#managedProcess || this.#managedProcess.exitCode !== null) {
const port = this.#baseUrl.port || (this.#baseUrl.protocol === 'https:' ? '443' : '80');
this.#managedProcess = spawn(this.#dshBin, [
'web', '--host', this.#baseUrl.hostname, '--port', port,
], {
cwd: this.#workspace,
env: process.env,
stdio: ['ignore', 'inherit', 'inherit'],
});
this.#managedProcess.on('error', (error) => {
console.error('[dsh-weixin] failed to start Harness:', error.message);
});
}
const deadline = Date.now() + 60_000;
let lastError;
while (Date.now() < deadline) {
await sleep(1_000);
try {
return await this.health();
} catch (error) {
lastError = error;
}
}
throw new Error(`Harness did not become ready: ${lastError?.message ?? 'timeout'}`);
}
async listWorkspaces(options = {}) {
await this.ensureRunning();
return workspacePaths(await this.rpc('workspace.list', {}, 30_000, options));
}
async listWorkspaceSessions(workspacePath, options = {}) {
await this.ensureRunning();
const workspaceList = await this.rpc('workspace.list', {}, 30_000, options);
const workspace = workspaceFromList(workspacePath, workspaceList);
if (!workspace) return { workspace: workspacePath, sessions: [] };
const sessionList = await this.rpc('session.list', {}, 30_000, options);
return workspaceSessions(workspace, workspaceList.archivedSessionIds, sessionList);
}
async adoptWorkspaceSession(value, options = {}) {
return adoptRegisteredWorkspaceSession(this, value, options);
}
async workspaceId(options = {}) {
const workspace = options.workspace ?? this.#workspace;
const { items } = await this.rpc('workspace.list', {});
const existing = items.find((item) => item.path === workspace);
if (existing) return existing.workspaceId;
const created = await this.rpc('workspace.create', { path: workspace });
return created.workspace.workspaceId;
}
async createSession(options = {}) {
await this.ensureRunning();
const workspaceId = await this.workspaceId(options);
const created = await this.rpc('session.create', {
workspaceId,
agentPreset: this.#agentPreset,
});
return created.sessionId;
}
async sessionExists(sessionId) {
try {
await this.rpc('session.history', { sessionId, maxMessages: 1 });
return true;
} catch (error) {
if (error instanceof HarnessRpcError && error.code === 'session-not-found') return false;
throw error;
}
}
async ask(sessionId, text, options = {}) {
if (typeof options === 'number') options = { timeoutMs: options };
const timeoutMs = options.timeoutMs ?? 600_000;
const onUpdate = typeof options.onUpdate === 'function' ? options.onUpdate : null;
await this.ensureRunning();
const before = await this.rpc('session.history', { sessionId, maxMessages: 1 });
const baselineSeq = Math.max(-1, ...(before.events ?? []).map(({ event }) => event.seq ?? -1));
const promptRpcId = `weixin-${randomUUID()}`;
const tracker = new HarnessReplyTracker({ promptRpcId, afterSeq: baselineSeq });
await this.rpc('session.prompt', {
sessionId,
mode: 'queue',
content: [{ type: 'text', text }],
clientTimeZone: Intl.DateTimeFormat().resolvedOptions().timeZone,
}, 30_000, { rpcId: promptRpcId });
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
await sleep(300);
const history = await this.rpc('session.history', { sessionId, maxMessages: 50 });
const update = tracker.consume(history.events ?? []);
if (update && onUpdate) {
try {
await onUpdate(update);
} catch (error) {
console.warn('[dsh-weixin] ignored a progress update failure:', error.message);
}
}
if (!tracker.finished) continue;
if (tracker.answer) return tracker.answer;
throw new Error(
`Harness turn ended without a text reply${tracker.reason ? ` (${JSON.stringify(tracker.reason)})` : ''}`,
);
}
throw new Error(`Harness reply timed out after ${Math.round(timeoutMs / 1_000)} seconds`);
}
stopManagedProcess() {
if (this.#managedProcess?.exitCode === null) this.#managedProcess.kill('SIGTERM');
}
}

View file

@ -3,9 +3,16 @@ import {
splitWeixinText,
weixinMessageId,
} from './weixin-api.mjs';
import {
harnessAnswerForQuestion,
harnessQuestionText,
validHarnessQuestion,
} from '../shared/harness-question.mjs';
import { runWorkspaceCommand } from '../shared/workspace-command.mjs';
import { askInWorkspaceSession } from '../shared/workspace-session.mjs';
const INTERACTION_RESOLVED_TEXT = '这个问题已在其他客户端处理,无需再次回答。';
const HELP_TEXT = [
'微信已连接 DeepSeek Harness。',
'',
@ -23,6 +30,16 @@ function conversationKey(userId) {
return `p2p:${userId}`;
}
function nonEmptyString(value) {
return typeof value === 'string' && value.trim() ? value.trim() : null;
}
function canClaimInteractionReply(message, pending) {
return pending.questions[pending.index]
&& nonEmptyString(message?.from_user_id) === pending.actor
&& nonEmptyString(extractWeixinText(message));
}
export function createWeixinBridgeStatus() {
return {
messagesReceived: 0,
@ -46,7 +63,11 @@ export class WeixinHarnessBridge {
#logger;
#replyTimeoutMs;
#maxMessageChars;
#signal;
#queues = new Map();
#pendingInteractions = new Map();
#interactionKeys = new Map();
#acceptedMessageIds = new Set();
constructor({
api,
@ -59,6 +80,7 @@ export class WeixinHarnessBridge {
logger = console,
replyTimeoutMs = 600_000,
maxMessageChars = 4_000,
signal,
}) {
if (!api || typeof api.sendText !== 'function') throw new TypeError('Weixin API is required');
if (!baseUrl || !token || !ownerUserId) throw new TypeError('Weixin account credentials are required');
@ -73,6 +95,7 @@ export class WeixinHarnessBridge {
this.#logger = logger;
this.#replyTimeoutMs = replyTimeoutMs;
this.#maxMessageChars = maxMessageChars;
this.#signal = signal;
}
get status() {
@ -80,31 +103,73 @@ export class WeixinHarnessBridge {
}
accept(message) {
const sender = typeof message?.from_user_id === 'string' ? message.from_user_id : '';
const previous = this.#queues.get(sender) ?? Promise.resolve();
if (this.#signal?.aborted) return Promise.resolve();
if (message?.message_type === 2) return Promise.resolve();
const messageId = weixinMessageId(message);
const sender = nonEmptyString(message?.from_user_id);
if (!messageId || !sender || this.#state.hasSeen(messageId)
|| this.#acceptedMessageIds.has(messageId)) return Promise.resolve();
this.#acceptedMessageIds.add(messageId);
const key = conversationKey(sender);
const pending = this.#pendingInteractions.get(key);
if (pending?.submitting || pending?.claimedReplyMessageId) {
return this.#enqueueMessage(message, messageId, key);
}
if (pending) {
if (canClaimInteractionReply(message, pending)) {
pending.claimedReplyMessageId = messageId;
}
const previous = pending.queue ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#processInteractionReply(message, messageId, key, pending))
.catch((error) => this.#handleInteractionFailure(message, messageId, error))
.finally(() => {
this.#acceptedMessageIds.delete(messageId);
if (pending.claimedReplyMessageId === messageId) pending.claimedReplyMessageId = null;
if (pending.queue === current) pending.queue = null;
});
pending.queue = current;
return current;
}
return this.#enqueueMessage(message, messageId, key);
}
#enqueueMessage(message, messageId, key, {
releaseMessageId = true,
alreadyRecorded = false,
} = {}) {
const previous = this.#queues.get(key) ?? Promise.resolve();
const current = previous
.catch(() => undefined)
.then(() => this.#process(message))
.then(() => this.#process(message, key, { alreadyRecorded }))
.finally(() => {
if (this.#queues.get(sender) === current) this.#queues.delete(sender);
if (releaseMessageId) this.#acceptedMessageIds.delete(messageId);
if (this.#queues.get(key) === current) this.#queues.delete(key);
});
this.#queues.set(sender, current);
this.#queues.set(key, current);
return current;
}
async waitForIdle() {
await Promise.allSettled([...this.#queues.values()]);
await Promise.allSettled([
...this.#queues.values(),
...[...this.#pendingInteractions.values()].flatMap((pending) => (
pending.queue ? [pending.queue] : []
)),
]);
}
async #process(message) {
if (message?.message_type === 2) return;
async #process(message, key, { alreadyRecorded = false } = {}) {
this.#signal?.throwIfAborted();
const messageId = weixinMessageId(message);
const sender = typeof message?.from_user_id === 'string' ? message.from_user_id : '';
const sender = nonEmptyString(message?.from_user_id);
if (!messageId || !sender) return;
if (this.#state.hasSeen(messageId)) return;
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
if (!alreadyRecorded) {
if (this.#state.hasSeen(messageId)) return;
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
}
if (sender !== this.#ownerUserId) {
this.#status.messagesRejected += 1;
this.#status.lastRejectedAt = new Date().toISOString();
@ -122,14 +187,13 @@ export class WeixinHarnessBridge {
}
const command = text.trim().toLowerCase();
const key = conversationKey(sender);
if (command === '/help') {
await this.#send(sender, HELP_TEXT, contextToken, runId);
await this.#state.markSeen(messageId);
return;
}
if (command === '/status') {
await this.#harness.ensureRunning();
await this.#harness.ensureRunning({ signal: this.#signal });
await this.#send(sender, '微信与 DeepSeek Harness 连接正常。', contextToken, runId);
await this.#state.markSeen(messageId);
return;
@ -149,19 +213,37 @@ export class WeixinHarnessBridge {
return;
}
const { answer } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
text,
askOptions: { timeoutMs: this.#replyTimeoutMs },
});
let answer;
try {
({ answer } = await askInWorkspaceSession({
harness: this.#harness,
state: this.#state,
key,
text,
createOptions: { signal: this.#signal },
existsOptions: { signal: this.#signal },
askOptions: {
timeoutMs: this.#replyTimeoutMs,
signal: this.#signal,
onInteraction: (interaction) => this.#handleInteraction(interaction, {
key,
actor: sender,
contextToken,
runId,
}),
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
},
}));
} finally {
await this.#cancelPendingInteraction(key);
}
await this.#send(sender, answer, contextToken, runId);
await this.#state.markSeen(messageId);
this.#status.messagesReplied += 1;
this.#status.lastReplyAt = new Date().toISOString();
this.#status.lastError = null;
} catch (error) {
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-weixin] failed to process an inbound message:', error);
try {
@ -173,6 +255,304 @@ export class WeixinHarnessBridge {
}
}
async #processInteractionReply(message, messageId, key, expected) {
this.#signal?.throwIfAborted();
const current = this.#pendingInteractions.get(key);
const claimed = expected.claimedReplyMessageId === messageId;
if (!current || current !== expected || current.submitting) {
if (claimed && (!current || current !== expected)) {
return this.#discardResolvedInteractionReply(message, messageId);
}
return this.#enqueueMessage(message, messageId, key, { releaseMessageId: false });
}
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
const text = nonEmptyString(extractWeixinText(message));
const contextToken = nonEmptyString(message?.context_token) ?? undefined;
const runId = nonEmptyString(message?.run_id) ?? undefined;
if (!text) {
await this.#send(
expected.actor,
'请用文字回答当前问题。',
contextToken,
runId,
);
return;
}
const pending = this.#pendingInteractions.get(key);
if (!pending || pending !== expected || pending.submitting) {
if (claimed && (!pending || pending !== expected)) {
await this.#send(
expected.actor,
INTERACTION_RESOLVED_TEXT,
contextToken,
runId,
);
return;
}
return this.#enqueueMessage(message, messageId, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
pending.contextToken = contextToken;
pending.runId = runId;
if (pending.needsPresentation) {
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = '微信交互问题发送失败。';
this.#logger.error?.('[dsh-weixin] failed to retry an interaction question');
pending.interaction.reconnect?.();
return;
}
const presentedPending = this.#pendingInteractions.get(key);
if (!presentedPending || presentedPending !== expected || presentedPending.submitting) {
if (claimed && (!presentedPending || presentedPending !== expected)) {
await this.#send(
expected.actor,
INTERACTION_RESOLVED_TEXT,
contextToken,
runId,
).catch(() => undefined);
return;
}
return this.#enqueueMessage(message, messageId, key, {
releaseMessageId: false,
alreadyRecorded: true,
});
}
}
const question = pending.questions[pending.index];
if (!question) return;
pending.answers.push(harnessAnswerForQuestion(question, text));
pending.index += 1;
if (pending.index < pending.questions.length) {
if (pending.claimedReplyMessageId === messageId) {
pending.claimedReplyMessageId = null;
}
pending.needsPresentation = true;
try {
await this.#presentInteraction(pending);
} catch {
this.#status.lastError = '微信交互问题发送失败。';
this.#logger.error?.('[dsh-weixin] failed to send the next interaction question');
pending.interaction.reconnect?.();
}
return;
}
pending.submitting = true;
try {
await pending.interaction.respond({
ok: true,
value: {
sessionId: pending.sessionId,
answer: { answers: pending.answers },
},
});
this.#clearPendingInteraction(key, pending.interactionId);
this.#status.lastError = null;
} catch (error) {
if (this.#signal?.aborted) return;
if (error?.code === 'interaction-not-pending') {
this.#clearPendingInteraction(key, pending.interactionId);
await this.#send(
pending.actor,
INTERACTION_RESOLVED_TEXT,
pending.contextToken,
pending.runId,
).catch(() => undefined);
return;
}
if (this.#pendingInteractions.get(key) !== pending) return;
pending.submitting = false;
pending.answers.pop();
pending.index -= 1;
this.#status.lastError = '回答提交失败。';
this.#logger.error?.('[dsh-weixin] failed to answer a Harness interaction');
await this.#send(
pending.actor,
'回答提交失败,请重新发送当前问题的答案。',
pending.contextToken,
pending.runId,
).catch(() => undefined);
}
}
async #handleInteraction(interaction, {
key,
actor,
contextToken,
runId,
}) {
// Approval remains fail-closed until #5 adds an authenticated policy.
if (interaction?.kind !== 'question') return;
const questions = interaction?.payload?.questions;
const interactionId = typeof interaction?.interactionId === 'string'
? interaction.interactionId
: interaction?.rpcId;
if (typeof interaction?.rpcId !== 'string'
|| typeof interactionId !== 'string'
|| typeof interaction.sessionId !== 'string'
|| !Array.isArray(questions)
|| questions.length === 0
|| questions.some((question) => !validHarnessQuestion(question))) {
this.#logger.warn?.('[dsh-weixin] ignored an invalid Harness question interaction');
return;
}
if (interaction.recovered === true) {
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'Weixin safely cancelled an interaction left by an earlier client.',
details: {},
},
});
await this.#send(
actor,
'检测到这个 Session 中遗留的待回答问题,已安全取消并继续处理你刚才的消息。',
contextToken,
runId,
).catch(() => undefined);
return;
}
const existing = this.#pendingInteractions.get(key);
if (existing?.interactionId === interactionId) {
existing.interaction = interaction;
if (existing.needsPresentation) await this.#presentInteraction(existing);
return;
}
if (this.#interactionKeys.has(interactionId)) return;
if (existing) {
await interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'Weixin is already handling another user interaction.',
details: {},
},
});
return;
}
const pending = {
kind: 'question',
interactionId,
sessionId: interaction.sessionId,
interaction,
actor,
questions,
answers: [],
index: 0,
contextToken,
runId,
queue: null,
claimedReplyMessageId: null,
presentationPromise: null,
submitting: false,
needsPresentation: true,
};
this.#pendingInteractions.set(key, pending);
this.#interactionKeys.set(interactionId, key);
await this.#presentInteraction(pending);
}
#handleInteractionResolved(resolution) {
const interactionId = resolution?.interactionId;
if (resolution?.kind !== 'question' || typeof interactionId !== 'string') return;
const key = this.#interactionKeys.get(interactionId);
if (!key) return;
this.#clearPendingInteraction(key, interactionId);
}
#presentInteraction(pending) {
if (!pending.needsPresentation) return Promise.resolve();
if (pending.presentationPromise) return pending.presentationPromise;
const question = pending.questions[pending.index];
if (!question) return Promise.resolve();
const presentation = this.#send(
pending.actor,
harnessQuestionText(question, pending.index, pending.questions.length),
pending.contextToken,
pending.runId,
).then(() => {
pending.needsPresentation = false;
}).finally(() => {
if (pending.presentationPromise === presentation) pending.presentationPromise = null;
});
pending.presentationPromise = presentation;
return presentation;
}
async #discardResolvedInteractionReply(message, messageId) {
if (this.#state.hasSeen(messageId)) return;
await this.#state.markSeen(messageId);
this.#status.messagesReceived += 1;
this.#status.lastMessageAt = new Date().toISOString();
await this.#send(
nonEmptyString(message?.from_user_id),
INTERACTION_RESOLVED_TEXT,
nonEmptyString(message?.context_token) ?? undefined,
nonEmptyString(message?.run_id) ?? undefined,
).catch(() => undefined);
}
#takePendingInteraction(key, interactionId) {
const pending = this.#pendingInteractions.get(key);
if (!pending
|| (interactionId !== undefined && pending.interactionId !== interactionId)) return null;
this.#pendingInteractions.delete(key);
this.#interactionKeys.delete(pending.interactionId);
return pending;
}
#clearPendingInteraction(key, interactionId) {
return this.#takePendingInteraction(key, interactionId) !== null;
}
async #cancelPendingInteraction(key) {
const pending = this.#takePendingInteraction(key);
if (!pending || pending.kind !== 'question') return;
try {
await pending.interaction.respond({
ok: false,
error: {
code: 'cancelled',
message: 'The Weixin interaction ended before the user answered.',
details: {},
},
}, { signal: AbortSignal.timeout(5_000) });
} catch (error) {
if (error?.code !== 'interaction-not-pending') {
this.#logger.warn?.('[dsh-weixin] failed to cancel a pending Harness interaction');
}
}
}
async #handleInteractionFailure(message, messageId, error) {
if (this.#signal?.aborted) return;
this.#status.lastError = error?.message ?? String(error);
this.#logger.error?.('[dsh-weixin] failed to process an interaction reply:', error);
if (!this.#state.hasSeen(messageId)) {
await this.#state.markSeen(messageId).catch(() => undefined);
}
await this.#send(
nonEmptyString(message?.from_user_id),
'消息处理失败,请稍后重试。',
nonEmptyString(message?.context_token) ?? undefined,
nonEmptyString(message?.run_id) ?? undefined,
).catch(() => undefined);
}
async #send(toUserId, text, contextToken, runId) {
for (const chunk of splitWeixinText(text, this.#maxMessageChars)) {
await this.#api.sendText({

View file

@ -116,6 +116,8 @@ export class WeixinRuntime {
await this.#harness.ensureRunning();
this.#status.harnessReachable = true;
await this.#notifyStart();
this.#abortController = new AbortController();
const signal = this.#abortController.signal;
this.#bridge = new WeixinHarnessBridge({
api: this.#api,
baseUrl: this.#config.baseUrl,
@ -127,12 +129,11 @@ export class WeixinRuntime {
logger: this.#logger,
replyTimeoutMs: this.#replyTimeoutMs,
maxMessageChars: this.#maxMessageChars,
signal,
});
this.#abortController = new AbortController();
this.#status.ready = true;
this.#status.weixinConnectionState = 'connected';
this.#status.lastCheckedAt = Date.now();
const signal = this.#abortController.signal;
this.#monitor = this.#runMonitor(signal).catch((error) => {
if (signal.aborted) return;
this.#status.ready = false;
@ -142,6 +143,9 @@ export class WeixinRuntime {
});
return this.status;
} catch (error) {
this.#abortController?.abort();
this.#abortController = null;
this.#bridge = null;
this.#status.ready = false;
this.#status.weixinConnectionState = 'failed';
this.#status.lastError = error?.message ?? String(error);
@ -195,7 +199,13 @@ export class WeixinRuntime {
this.#status.lastError = null;
for (const message of response?.msgs ?? []) {
await this.#bridge.accept(message);
void this.#bridge.accept(message).catch((error) => {
if (signal.aborted) return;
this.#logger.error?.(
`[dsh-weixin] account ${this.#config.botId} message handling failed:`,
error,
);
});
}
if (typeof response?.get_updates_buf === 'string' && response.get_updates_buf) {
await this.#state.setGetUpdatesBuf(response.get_updates_buf);

View file

@ -1,3 +1,11 @@
import { HarnessClient } from '../weixin/harness-client.mjs';
import { HarnessClient } from '../shared/harness-client.mjs';
export class WhatsappHarnessClient extends HarnessClient {}
export class WhatsappHarnessClient extends HarnessClient {
constructor(options) {
super({
...options,
rpcIdPrefix: 'whatsapp',
logPrefix: 'dsh-whatsapp',
});
}
}

View file

@ -256,6 +256,7 @@ export class WhatsappRuntime {
status: this.#status,
logger: this.#logger,
replyTimeoutMs: this.#replyTimeoutMs,
signal: controller.signal,
});
const now = Date.now();
this.#status.ready = true;