mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-10 22:03:18 +08:00
Add WeCom QR bot channel
This commit is contained in:
parent
601271a7f8
commit
4874d3e354
33 changed files with 5200 additions and 440 deletions
170
src/channels/wecom/config-store.mjs
Normal file
170
src/channels/wecom/config-store.mjs
Normal file
|
|
@ -0,0 +1,170 @@
|
|||
import { createHash } from 'node:crypto';
|
||||
import { mkdir, readFile, rename, unlink, writeFile } from 'node:fs/promises';
|
||||
import { dirname } from 'node:path';
|
||||
|
||||
const EMPTY_DOCUMENT = Object.freeze({ version: 1, bots: Object.freeze([]) });
|
||||
|
||||
function cleanString(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
function safeIntegrationId(value) {
|
||||
const id = cleanString(value);
|
||||
return id && /^wecom_[a-f0-9]{24}$/.test(id) ? id : null;
|
||||
}
|
||||
|
||||
function safeSecretRef(value) {
|
||||
const ref = cleanString(value);
|
||||
return ref && /^DSH_WECOM_BOT_SECRET_[A-F0-9]{24}$/.test(ref) ? ref : null;
|
||||
}
|
||||
|
||||
export function deriveWecomBotIdentity(remoteBotId) {
|
||||
const raw = cleanString(remoteBotId);
|
||||
if (!raw) throw new TypeError('Enterprise WeChat bot ID is required');
|
||||
const digest = createHash('sha256').update(raw).digest('hex').slice(0, 24);
|
||||
return {
|
||||
botId: `wecom_${digest}`,
|
||||
secretRef: `DSH_WECOM_BOT_SECRET_${digest.toUpperCase()}`,
|
||||
};
|
||||
}
|
||||
|
||||
export function maskWecomBotId(remoteBotId) {
|
||||
const value = cleanString(remoteBotId) ?? '';
|
||||
if (!value) return '企业微信机器人';
|
||||
if (value.length <= 10) return `${value.slice(0, 3)}•••`;
|
||||
return `${value.slice(0, 6)}••••${value.slice(-4)}`;
|
||||
}
|
||||
|
||||
function normalizeBot(value) {
|
||||
if (!value || typeof value !== 'object') return null;
|
||||
const botId = safeIntegrationId(value.botId);
|
||||
const remoteBotId = cleanString(value.remoteBotId);
|
||||
const secretRef = safeSecretRef(value.secretRef);
|
||||
if (!botId || !remoteBotId || !secretRef) return null;
|
||||
const derived = deriveWecomBotIdentity(remoteBotId);
|
||||
if (derived.botId !== botId || derived.secretRef !== secretRef) return null;
|
||||
return Object.freeze({
|
||||
botId,
|
||||
remoteBotId,
|
||||
secretRef,
|
||||
createdAt: cleanString(value.createdAt) ?? new Date().toISOString(),
|
||||
connectedAt: cleanString(value.connectedAt),
|
||||
});
|
||||
}
|
||||
|
||||
function normalizeDocument(value) {
|
||||
if (!value || value.version !== 1 || !Array.isArray(value.bots)) return null;
|
||||
const bots = value.bots.map(normalizeBot);
|
||||
if (bots.some((bot) => bot === null)) return null;
|
||||
const ids = new Set();
|
||||
const remoteIds = new Set();
|
||||
const refs = new Set();
|
||||
for (const bot of bots) {
|
||||
if (ids.has(bot.botId) || remoteIds.has(bot.remoteBotId) || refs.has(bot.secretRef)) return null;
|
||||
ids.add(bot.botId);
|
||||
remoteIds.add(bot.remoteBotId);
|
||||
refs.add(bot.secretRef);
|
||||
}
|
||||
return Object.freeze({ version: 1, bots: Object.freeze(bots) });
|
||||
}
|
||||
|
||||
export class WecomConfigStore {
|
||||
#path;
|
||||
#value = EMPTY_DOCUMENT;
|
||||
#writeQueue = Promise.resolve();
|
||||
|
||||
constructor(path) {
|
||||
this.#path = path;
|
||||
}
|
||||
|
||||
async load() {
|
||||
try {
|
||||
const normalized = normalizeDocument(JSON.parse(await readFile(this.#path, 'utf8')));
|
||||
if (!normalized) throw new Error('dsh-im Enterprise WeChat config contains invalid bot data');
|
||||
this.#value = normalized;
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
this.#value = EMPTY_DOCUMENT;
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
list() {
|
||||
return structuredClone(this.#value.bots);
|
||||
}
|
||||
|
||||
get(botId) {
|
||||
const bot = this.#value.bots.find((candidate) => candidate.botId === botId);
|
||||
return bot ? structuredClone(bot) : null;
|
||||
}
|
||||
|
||||
getByRemoteBotId(remoteBotId) {
|
||||
const bot = this.#value.bots.find((candidate) => candidate.remoteBotId === remoteBotId);
|
||||
return bot ? structuredClone(bot) : null;
|
||||
}
|
||||
|
||||
async save(value) {
|
||||
const normalized = normalizeBot(value);
|
||||
if (!normalized) throw new Error('Refusing to persist incomplete Enterprise WeChat bot data');
|
||||
return this.#mutate((bots) => {
|
||||
const remoteCollision = bots.find(
|
||||
(bot) => bot.remoteBotId === normalized.remoteBotId && bot.botId !== normalized.botId,
|
||||
);
|
||||
const refCollision = bots.find(
|
||||
(bot) => bot.secretRef === normalized.secretRef && bot.botId !== normalized.botId,
|
||||
);
|
||||
if (remoteCollision || refCollision) throw new Error('Duplicate Enterprise WeChat bot identity');
|
||||
const index = bots.findIndex((bot) => bot.botId === normalized.botId);
|
||||
if (index === -1) bots.push(normalized);
|
||||
else bots[index] = normalized;
|
||||
return structuredClone(normalized);
|
||||
});
|
||||
}
|
||||
|
||||
async remove(botId) {
|
||||
if (!safeIntegrationId(botId)) throw new TypeError('Invalid Enterprise WeChat bot ID');
|
||||
return this.#mutate((bots) => {
|
||||
const index = bots.findIndex((bot) => bot.botId === botId);
|
||||
if (index === -1) return null;
|
||||
const [removed] = bots.splice(index, 1);
|
||||
return structuredClone(removed);
|
||||
});
|
||||
}
|
||||
|
||||
async clear() {
|
||||
const operation = this.#writeQueue.then(async () => {
|
||||
try {
|
||||
await unlink(this.#path);
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
}
|
||||
this.#value = EMPTY_DOCUMENT;
|
||||
});
|
||||
this.#writeQueue = operation.then(() => undefined, () => undefined);
|
||||
await operation;
|
||||
}
|
||||
|
||||
async #mutate(mutator) {
|
||||
let result;
|
||||
const operation = this.#writeQueue.then(async () => {
|
||||
const bots = [...this.#value.bots];
|
||||
result = mutator(bots);
|
||||
const document = Object.freeze({ version: 1, bots: Object.freeze(bots) });
|
||||
await this.#write(document);
|
||||
this.#value = document;
|
||||
});
|
||||
this.#writeQueue = operation.then(() => undefined, () => undefined);
|
||||
await operation;
|
||||
return result;
|
||||
}
|
||||
|
||||
async #write(document) {
|
||||
await mkdir(dirname(this.#path), { recursive: true, mode: 0o700 });
|
||||
const temporary = `${this.#path}.tmp`;
|
||||
await writeFile(temporary, `${JSON.stringify(document, null, 2)}\n`, {
|
||||
encoding: 'utf8',
|
||||
mode: 0o600,
|
||||
});
|
||||
await rename(temporary, this.#path);
|
||||
}
|
||||
}
|
||||
3
src/channels/wecom/harness-client.mjs
Normal file
3
src/channels/wecom/harness-client.mjs
Normal file
|
|
@ -0,0 +1,3 @@
|
|||
import { HarnessClient } from '../weixin/harness-client.mjs';
|
||||
|
||||
export class WecomHarnessClient extends HarnessClient {}
|
||||
95
src/channels/wecom/qr-auth.mjs
Normal file
95
src/channels/wecom/qr-auth.mjs
Normal file
|
|
@ -0,0 +1,95 @@
|
|||
const GENERATE_URL = 'https://work.weixin.qq.com/ai/qc/generate';
|
||||
const POLL_URL = 'https://work.weixin.qq.com/ai/qc/query_result';
|
||||
const QR_TTL_MS = 5 * 60_000;
|
||||
const POLL_INTERVAL_MS = 3_000;
|
||||
|
||||
function defaultPlatform() {
|
||||
if (process.platform === 'win32') return 2;
|
||||
if (process.platform === 'linux') return 3;
|
||||
return 1;
|
||||
}
|
||||
|
||||
function cleanString(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
function safeVerificationUrl(value) {
|
||||
const raw = cleanString(value);
|
||||
if (!raw) return null;
|
||||
try {
|
||||
const url = new URL(raw);
|
||||
return url.protocol === 'https:' && url.hostname === 'work.weixin.qq.com' && (!url.port || url.port === '443')
|
||||
? url.href : null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function combinedSignal(signal, timeoutMs) {
|
||||
return signal ? AbortSignal.any([signal, AbortSignal.timeout(timeoutMs)]) : AbortSignal.timeout(timeoutMs);
|
||||
}
|
||||
|
||||
async function requestJson(fetchImpl, url, signal) {
|
||||
const response = await fetchImpl(url, {
|
||||
method: 'GET',
|
||||
redirect: 'error',
|
||||
signal: combinedSignal(signal, 10_000),
|
||||
headers: { accept: 'application/json' },
|
||||
});
|
||||
if (!response.ok) throw new Error(`Enterprise WeChat QR service returned HTTP ${response.status}`);
|
||||
return response.json();
|
||||
}
|
||||
|
||||
export class WecomQrAuth {
|
||||
#fetch;
|
||||
#clock;
|
||||
#source;
|
||||
#platform;
|
||||
|
||||
constructor({
|
||||
fetch: fetchImpl = globalThis.fetch,
|
||||
clock = () => Date.now(),
|
||||
source = 'deepseek-harness',
|
||||
platform = defaultPlatform(),
|
||||
} = {}) {
|
||||
if (typeof fetchImpl !== 'function') throw new TypeError('A fetch implementation is required');
|
||||
this.#fetch = fetchImpl;
|
||||
this.#clock = clock;
|
||||
this.#source = source;
|
||||
this.#platform = [1, 2, 3].includes(platform) ? platform : defaultPlatform();
|
||||
}
|
||||
|
||||
async start({ signal } = {}) {
|
||||
const url = new URL(GENERATE_URL);
|
||||
url.searchParams.set('source', this.#source);
|
||||
url.searchParams.set('plat', String(this.#platform));
|
||||
const body = await requestJson(this.#fetch, url, signal);
|
||||
const scode = cleanString(body?.data?.scode);
|
||||
const verificationUrl = safeVerificationUrl(body?.data?.auth_url);
|
||||
if (!scode || !verificationUrl) throw new Error('Enterprise WeChat QR service returned invalid data');
|
||||
return {
|
||||
scode,
|
||||
verificationUrl,
|
||||
expiresAt: this.#clock() + QR_TTL_MS,
|
||||
pollIntervalMs: POLL_INTERVAL_MS,
|
||||
};
|
||||
}
|
||||
|
||||
async poll({ scode, signal } = {}) {
|
||||
const code = cleanString(scode);
|
||||
if (!code) throw new TypeError('Enterprise WeChat QR poll code is required');
|
||||
const url = new URL(POLL_URL);
|
||||
url.searchParams.set('scode', code);
|
||||
const body = await requestJson(this.#fetch, url, signal);
|
||||
const state = cleanString(body?.data?.status)?.toLowerCase();
|
||||
if (state === 'success') {
|
||||
const remoteBotId = cleanString(body?.data?.bot_info?.botid);
|
||||
const secret = cleanString(body?.data?.bot_info?.secret);
|
||||
if (!remoteBotId || !secret) throw new Error('Enterprise WeChat QR result omitted bot credentials');
|
||||
return { status: 'success', remoteBotId, secret };
|
||||
}
|
||||
if (['expired', 'timeout'].includes(state)) return { status: 'expired' };
|
||||
if (['fail', 'failed', 'error'].includes(state)) return { status: 'failed' };
|
||||
return { status: 'waiting' };
|
||||
}
|
||||
}
|
||||
90
src/channels/wecom/state-store.mjs
Normal file
90
src/channels/wecom/state-store.mjs
Normal file
|
|
@ -0,0 +1,90 @@
|
|||
import { mkdir, readFile, rename, unlink, writeFile } from 'node:fs/promises';
|
||||
import { dirname } from 'node:path';
|
||||
|
||||
const EMPTY_STATE = Object.freeze({ version: 1, sessions: {}, seenMessageIds: [] });
|
||||
|
||||
function normalizeState(value) {
|
||||
if (!value || typeof value !== 'object') return structuredClone(EMPTY_STATE);
|
||||
const sessions = {};
|
||||
if (value.sessions && typeof value.sessions === 'object' && !Array.isArray(value.sessions)) {
|
||||
for (const [key, sessionId] of Object.entries(value.sessions)) {
|
||||
if (typeof key === 'string' && typeof sessionId === 'string' && sessionId) sessions[key] = sessionId;
|
||||
}
|
||||
}
|
||||
return {
|
||||
version: 1,
|
||||
sessions,
|
||||
seenMessageIds: Array.isArray(value.seenMessageIds)
|
||||
? value.seenMessageIds.filter((id) => typeof id === 'string').slice(-1_000)
|
||||
: [],
|
||||
};
|
||||
}
|
||||
|
||||
export class WecomStateStore {
|
||||
#path;
|
||||
#state = structuredClone(EMPTY_STATE);
|
||||
#writeQueue = Promise.resolve();
|
||||
|
||||
constructor(path) {
|
||||
this.#path = path;
|
||||
}
|
||||
|
||||
async load() {
|
||||
try {
|
||||
this.#state = normalizeState(JSON.parse(await readFile(this.#path, 'utf8')));
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
this.#state = structuredClone(EMPTY_STATE);
|
||||
await this.#persist();
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
sessionFor(key) {
|
||||
return this.#state.sessions[key] ?? null;
|
||||
}
|
||||
|
||||
async setSession(key, sessionId) {
|
||||
this.#state.sessions[key] = sessionId;
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
async clearSession(key) {
|
||||
delete this.#state.sessions[key];
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
hasSeen(messageId) {
|
||||
return this.#state.seenMessageIds.includes(messageId);
|
||||
}
|
||||
|
||||
async markSeen(messageId) {
|
||||
if (this.hasSeen(messageId)) return;
|
||||
this.#state.seenMessageIds.push(messageId);
|
||||
if (this.#state.seenMessageIds.length > 1_000) {
|
||||
this.#state.seenMessageIds.splice(0, this.#state.seenMessageIds.length - 1_000);
|
||||
}
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
async remove() {
|
||||
try {
|
||||
await unlink(this.#path);
|
||||
} catch (error) {
|
||||
if (error?.code !== 'ENOENT') throw error;
|
||||
}
|
||||
this.#state = structuredClone(EMPTY_STATE);
|
||||
}
|
||||
|
||||
async #persist() {
|
||||
const snapshot = `${JSON.stringify(this.#state, null, 2)}\n`;
|
||||
const operation = this.#writeQueue.then(async () => {
|
||||
await mkdir(dirname(this.#path), { recursive: true, mode: 0o700 });
|
||||
const temporary = `${this.#path}.tmp`;
|
||||
await writeFile(temporary, snapshot, { encoding: 'utf8', mode: 0o600 });
|
||||
await rename(temporary, this.#path);
|
||||
});
|
||||
this.#writeQueue = operation.then(() => undefined, () => undefined);
|
||||
await operation;
|
||||
}
|
||||
}
|
||||
243
src/channels/wecom/wecom-bridge.mjs
Normal file
243
src/channels/wecom/wecom-bridge.mjs
Normal file
|
|
@ -0,0 +1,243 @@
|
|||
import { generateReqId } from '@wecom/aibot-node-sdk';
|
||||
|
||||
const HELP_TEXT = [
|
||||
'企业微信机器人已连接 DeepSeek Harness。',
|
||||
'',
|
||||
'直接发送文字即可继续当前会话。',
|
||||
'/new 开启一个全新会话',
|
||||
'/status 检查连接状态',
|
||||
'/help 显示本帮助',
|
||||
].join('\n');
|
||||
const MAX_REPLY_BYTES = 18_000;
|
||||
|
||||
function bodyOf(frame) {
|
||||
return frame?.body && typeof frame.body === 'object' ? frame.body : {};
|
||||
}
|
||||
|
||||
function conversationKey(frame) {
|
||||
const body = bodyOf(frame);
|
||||
return body.chattype === 'group' ? `group:${body.chatid}` : `direct:${body.from?.userid}`;
|
||||
}
|
||||
|
||||
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
|
||||
.filter((item) => item?.msgtype === 'text' && typeof item.text?.content === 'string')
|
||||
.map((item) => item.text.content)
|
||||
.join('\n')
|
||||
.trim();
|
||||
}
|
||||
return '';
|
||||
}
|
||||
|
||||
function splitUtf8(text, maxBytes = MAX_REPLY_BYTES) {
|
||||
const source = String(text ?? '').trim();
|
||||
if (!source) return [];
|
||||
const chunks = [];
|
||||
let current = '';
|
||||
let bytes = 0;
|
||||
for (const character of source) {
|
||||
const size = Buffer.byteLength(character);
|
||||
if (current && bytes + size > maxBytes) {
|
||||
chunks.push(current);
|
||||
current = character;
|
||||
bytes = size;
|
||||
} else {
|
||||
current += character;
|
||||
bytes += size;
|
||||
}
|
||||
}
|
||||
if (current) chunks.push(current);
|
||||
return chunks;
|
||||
}
|
||||
|
||||
function progressText(update) {
|
||||
if (update?.type === 'text') return update.text;
|
||||
if (update?.type === 'tool') return `正在使用${update.name}…`;
|
||||
return update?.text;
|
||||
}
|
||||
|
||||
export function createWecomBridgeStatus() {
|
||||
return {
|
||||
messagesReceived: 0,
|
||||
messagesReplied: 0,
|
||||
messagesRejected: 0,
|
||||
lastMessageAt: null,
|
||||
lastReplyAt: null,
|
||||
lastRejectedAt: null,
|
||||
lastError: null,
|
||||
};
|
||||
}
|
||||
|
||||
export class WecomHarnessBridge {
|
||||
#client;
|
||||
#harness;
|
||||
#state;
|
||||
#status;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#generateReqId;
|
||||
#queues = new Map();
|
||||
|
||||
constructor({
|
||||
client,
|
||||
harness,
|
||||
state,
|
||||
status = createWecomBridgeStatus(),
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
generateStreamId = generateReqId,
|
||||
}) {
|
||||
if (!client || typeof client.replyStream !== 'function' || typeof client.sendMessage !== 'function') {
|
||||
throw new TypeError('Enterprise WeChat client is required');
|
||||
}
|
||||
if (!harness || !state) throw new TypeError('Harness client and state store are required');
|
||||
this.#client = client;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#status = status;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#generateReqId = generateStreamId;
|
||||
}
|
||||
|
||||
get status() {
|
||||
return structuredClone(this.#status);
|
||||
}
|
||||
|
||||
accept(frame) {
|
||||
const key = conversationKey(frame);
|
||||
const previous = this.#queues.get(key) ?? Promise.resolve();
|
||||
const current = previous
|
||||
.catch(() => undefined)
|
||||
.then(() => this.#process(frame))
|
||||
.finally(() => {
|
||||
if (this.#queues.get(key) === current) this.#queues.delete(key);
|
||||
});
|
||||
this.#queues.set(key, current);
|
||||
return current;
|
||||
}
|
||||
|
||||
async waitForIdle() {
|
||||
await Promise.allSettled([...this.#queues.values()]);
|
||||
}
|
||||
|
||||
async #sendActive(chatId, text) {
|
||||
for (const chunk of splitUtf8(text)) {
|
||||
await this.#client.sendMessage(chatId, { msgtype: 'markdown', markdown: { content: chunk } });
|
||||
}
|
||||
}
|
||||
|
||||
async #sendImmediate(frame, chatId, text) {
|
||||
const chunks = splitUtf8(text);
|
||||
if (chunks.length === 0) return;
|
||||
try {
|
||||
await this.#client.replyStream(frame, this.#generateReqId('stream'), chunks[0], true);
|
||||
for (const chunk of chunks.slice(1)) {
|
||||
await this.#client.sendMessage(chatId, { msgtype: 'markdown', markdown: { content: chunk } });
|
||||
}
|
||||
} catch {
|
||||
await this.#sendActive(chatId, text);
|
||||
}
|
||||
}
|
||||
|
||||
async #process(frame) {
|
||||
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();
|
||||
const text = messageText(frame);
|
||||
let streamId = null;
|
||||
let streamStarted = false;
|
||||
try {
|
||||
if (!text) {
|
||||
await this.#sendImmediate(frame, chatId, '目前支持文字、语音转写和图文混排中的文字消息。');
|
||||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
const command = text.toLowerCase();
|
||||
if (command === '/help') {
|
||||
await this.#sendImmediate(frame, chatId, HELP_TEXT);
|
||||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
if (command === '/status') {
|
||||
await this.#harness.ensureRunning();
|
||||
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, '已开启新会话。请发送你的问题。');
|
||||
await this.#state.markSeen(messageId);
|
||||
return;
|
||||
}
|
||||
|
||||
let sessionId = this.#state.sessionFor(key);
|
||||
if (!sessionId || !(await this.#harness.sessionExists(sessionId))) {
|
||||
sessionId = await this.#harness.createSession();
|
||||
await this.#state.setSession(key, sessionId);
|
||||
}
|
||||
|
||||
streamId = this.#generateReqId('stream');
|
||||
try {
|
||||
await this.#client.replyStream(frame, streamId, '正在思考中…', false);
|
||||
streamStarted = true;
|
||||
} catch (error) {
|
||||
this.#logger.warn?.('[dsh-im:wecom] unable to start a stream; using an active reply:', error);
|
||||
}
|
||||
|
||||
const answer = await this.#harness.ask(sessionId, text, {
|
||||
timeoutMs: this.#replyTimeoutMs,
|
||||
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,
|
||||
});
|
||||
|
||||
const chunks = splitUtf8(answer || '任务已完成,但没有生成可显示的文本。');
|
||||
let finalSent = false;
|
||||
if (streamStarted && chunks.length > 0) {
|
||||
try {
|
||||
await this.#client.replyStream(frame, streamId, chunks[0], true);
|
||||
for (const chunk of chunks.slice(1)) {
|
||||
await this.#client.sendMessage(chatId, { msgtype: 'markdown', markdown: { content: chunk } });
|
||||
}
|
||||
finalSent = true;
|
||||
} catch (error) {
|
||||
this.#logger.warn?.('[dsh-im:wecom] stream finalization failed; using an active reply:', error);
|
||||
}
|
||||
}
|
||||
if (!finalSent) await this.#sendActive(chatId, answer);
|
||||
await this.#state.markSeen(messageId);
|
||||
this.#status.messagesReplied += 1;
|
||||
this.#status.lastReplyAt = new Date().toISOString();
|
||||
this.#status.lastError = null;
|
||||
} catch (error) {
|
||||
this.#status.lastError = error?.message ?? String(error);
|
||||
this.#logger.error?.('[dsh-im:wecom] failed to process an inbound message');
|
||||
try {
|
||||
if (streamStarted && streamId) {
|
||||
await this.#client.replyStream(frame, streamId, '消息处理失败,请稍后重试。', true);
|
||||
} else {
|
||||
await this.#sendImmediate(frame, chatId, '消息处理失败,请稍后重试。');
|
||||
}
|
||||
await this.#state.markSeen(messageId);
|
||||
} catch {
|
||||
this.#logger.error?.('[dsh-im:wecom] failed to send the safe error reply');
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
427
src/channels/wecom/wecom-controller.mjs
Normal file
427
src/channels/wecom/wecom-controller.mjs
Normal file
|
|
@ -0,0 +1,427 @@
|
|||
import { randomUUID } from 'node:crypto';
|
||||
|
||||
import { deriveWecomBotIdentity, maskWecomBotId } from './config-store.mjs';
|
||||
|
||||
const ACTIVE_ATTEMPT_STATES = new Set(['pending', 'connecting']);
|
||||
const TERMINAL_ATTEMPT_STATES = new Set(['connected', 'failed', 'cancelled', 'expired']);
|
||||
|
||||
function cleanString(value) {
|
||||
return typeof value === 'string' && value.trim() ? value.trim() : null;
|
||||
}
|
||||
|
||||
function safeError(code, message) {
|
||||
return Object.freeze({ code, message });
|
||||
}
|
||||
|
||||
function publicAttempt(record) {
|
||||
if (!record) return null;
|
||||
return {
|
||||
attemptId: record.id,
|
||||
status: record.state,
|
||||
pollIntervalMs: record.pollIntervalMs,
|
||||
qrRevision: record.qrRevision,
|
||||
...(record.verificationUrl ? { verificationUrl: record.verificationUrl } : {}),
|
||||
...(record.expiresAt ? { expiresAt: record.expiresAt } : {}),
|
||||
...(record.botId ? { botId: record.botId } : {}),
|
||||
...(record.error ? { error: structuredClone(record.error) } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
export class WecomController {
|
||||
#qrAuth;
|
||||
#credentials;
|
||||
#configStore;
|
||||
#createRuntime;
|
||||
#deleteState;
|
||||
#logger;
|
||||
#runtimes = new Map();
|
||||
#errors = new Map();
|
||||
#attempts = new Map();
|
||||
#activeAttemptId = null;
|
||||
#transitions = new Map();
|
||||
#revision = 0;
|
||||
#closed = false;
|
||||
|
||||
constructor({
|
||||
qrAuth,
|
||||
credentials,
|
||||
configStore,
|
||||
createRuntime,
|
||||
deleteState = async () => {},
|
||||
logger = console,
|
||||
}) {
|
||||
if (!qrAuth || typeof qrAuth.start !== 'function' || typeof qrAuth.poll !== 'function') {
|
||||
throw new TypeError('Enterprise WeChat QR auth is required');
|
||||
}
|
||||
if (!credentials || typeof credentials.resolve !== 'function'
|
||||
|| typeof credentials.set !== 'function' || typeof credentials.unset !== 'function') {
|
||||
throw new TypeError('WecomController requires the DSH credential provider');
|
||||
}
|
||||
if (!configStore || typeof configStore.list !== 'function'
|
||||
|| typeof configStore.save !== 'function' || typeof configStore.remove !== 'function') {
|
||||
throw new TypeError('WecomController requires a config store');
|
||||
}
|
||||
if (typeof createRuntime !== 'function') throw new TypeError('createRuntime is required');
|
||||
this.#qrAuth = qrAuth;
|
||||
this.#credentials = credentials;
|
||||
this.#configStore = configStore;
|
||||
this.#createRuntime = createRuntime;
|
||||
this.#deleteState = deleteState;
|
||||
this.#logger = logger;
|
||||
}
|
||||
|
||||
async initialize() {
|
||||
if (this.#closed) return this.status();
|
||||
for (const config of this.#configStore.list()) {
|
||||
await this.#withBotTransition(config.botId, async () => {
|
||||
const existing = this.#runtimes.get(config.botId)?.status;
|
||||
if (this.#closed || existing?.ready || existing?.wecomConnectionState === 'connecting') return;
|
||||
const secret = await this.#resolveSecret(config.secretRef);
|
||||
if (!secret) {
|
||||
this.#errors.set(config.botId, safeError('missing-secret', '企业微信机器人凭据缺失,请移除后重新扫码。'));
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await this.#startRuntime(config, secret);
|
||||
this.#errors.delete(config.botId);
|
||||
} catch (error) {
|
||||
this.#errors.set(config.botId, safeError('connection-failed', '企业微信连接未就绪,插件会自动重试。'));
|
||||
this.#logger.warn?.(`[dsh-im:wecom] bot ${config.botId} failed to initialize`);
|
||||
} finally {
|
||||
this.#touch();
|
||||
}
|
||||
});
|
||||
}
|
||||
return this.status();
|
||||
}
|
||||
|
||||
async startProvisioning() {
|
||||
if (this.#closed) throw new Error('Enterprise WeChat controller is closed');
|
||||
if (this.#activeAttemptId) await this.cancelProvisioning(this.#activeAttemptId);
|
||||
const record = {
|
||||
id: randomUUID(),
|
||||
state: 'pending',
|
||||
createdAt: Date.now(),
|
||||
expiresAt: null,
|
||||
pollIntervalMs: 3_000,
|
||||
qrRevision: 1,
|
||||
verificationUrl: null,
|
||||
scode: null,
|
||||
botId: null,
|
||||
error: null,
|
||||
controller: new AbortController(),
|
||||
polling: null,
|
||||
transition: null,
|
||||
};
|
||||
this.#attempts.set(record.id, record);
|
||||
this.#activeAttemptId = record.id;
|
||||
this.#touch();
|
||||
try {
|
||||
const started = await this.#qrAuth.start({ signal: record.controller.signal });
|
||||
record.scode = cleanString(started.scode);
|
||||
record.verificationUrl = cleanString(started.verificationUrl);
|
||||
record.expiresAt = Number(started.expiresAt);
|
||||
record.pollIntervalMs = Math.min(10_000, Math.max(500, Number(started.pollIntervalMs) || 3_000));
|
||||
if (!record.scode || !record.verificationUrl || !Number.isFinite(record.expiresAt)) {
|
||||
throw new Error('Enterprise WeChat QR auth returned incomplete data');
|
||||
}
|
||||
this.#touch();
|
||||
return publicAttempt(record);
|
||||
} catch (error) {
|
||||
record.state = record.controller.signal.aborted ? 'cancelled' : 'failed';
|
||||
record.error = record.controller.signal.aborted
|
||||
? safeError('cancelled', '扫码绑定已取消。')
|
||||
: safeError('qr-start-failed', '无法生成企业微信二维码,请稍后重试。');
|
||||
this.#finishAttempt(record);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async registrationStatus(attemptId) {
|
||||
const record = this.#attempts.get(attemptId);
|
||||
if (!record || TERMINAL_ATTEMPT_STATES.has(record.state)) return publicAttempt(record);
|
||||
if (record.state === 'connecting') {
|
||||
await record.transition?.catch(() => undefined);
|
||||
return publicAttempt(record);
|
||||
}
|
||||
if (Date.now() >= record.expiresAt) {
|
||||
record.state = 'expired';
|
||||
record.error = safeError('expired', '企业微信二维码已过期,请重新生成。');
|
||||
record.controller.abort();
|
||||
this.#finishAttempt(record);
|
||||
return publicAttempt(record);
|
||||
}
|
||||
if (!record.polling) {
|
||||
const polling = this.#pollAttempt(record).finally(() => {
|
||||
if (record.polling === polling) record.polling = null;
|
||||
});
|
||||
record.polling = polling;
|
||||
}
|
||||
await record.polling.catch(() => undefined);
|
||||
return publicAttempt(record);
|
||||
}
|
||||
|
||||
async cancelProvisioning(attemptId) {
|
||||
const record = this.#attempts.get(attemptId);
|
||||
if (!record) return null;
|
||||
if (!TERMINAL_ATTEMPT_STATES.has(record.state)) {
|
||||
record.controller.abort();
|
||||
await Promise.allSettled([record.polling, record.transition].filter(Boolean));
|
||||
if (!TERMINAL_ATTEMPT_STATES.has(record.state)) record.state = 'cancelled';
|
||||
record.error ??= safeError('cancelled', '扫码绑定已取消。');
|
||||
this.#finishAttempt(record);
|
||||
}
|
||||
return publicAttempt(record);
|
||||
}
|
||||
|
||||
async reconnectBot(botId) {
|
||||
const config = this.#configStore.get(botId);
|
||||
if (!config) throw new Error('Unknown Enterprise WeChat bot');
|
||||
await this.#withBotTransition(botId, async () => {
|
||||
const secret = await this.#resolveSecret(config.secretRef);
|
||||
if (!secret) throw new Error('Enterprise WeChat bot secret is missing');
|
||||
try {
|
||||
await this.#startRuntime(config, secret);
|
||||
this.#errors.delete(botId);
|
||||
} catch (error) {
|
||||
this.#errors.set(botId, safeError('connection-failed', '企业微信连接仍未就绪,请稍后重试。'));
|
||||
throw error;
|
||||
} finally {
|
||||
this.#touch();
|
||||
}
|
||||
});
|
||||
return this.status();
|
||||
}
|
||||
|
||||
async deleteBot(botId) {
|
||||
const config = this.#configStore.get(botId);
|
||||
if (!config) throw new Error('Unknown Enterprise WeChat bot');
|
||||
await this.#withBotTransition(botId, async () => {
|
||||
const previous = await this.#credentials.resolve(config.secretRef).catch(() => undefined);
|
||||
await this.#stopRuntime(botId);
|
||||
try {
|
||||
await this.#credentials.unset(config.secretRef);
|
||||
await this.#configStore.remove(botId);
|
||||
} catch (error) {
|
||||
if (previous?.value) {
|
||||
await this.#credentials.set(config.secretRef, previous.value).catch(() => undefined);
|
||||
await this.#startRuntime(config, previous.value).catch(() => undefined);
|
||||
}
|
||||
throw new Error('Unable to remove the Enterprise WeChat bot safely.', { cause: error });
|
||||
}
|
||||
await this.#deleteState({ botId, config }).catch((error) => {
|
||||
this.#logger.warn?.(`[dsh-im:wecom] bot ${botId} state cleanup failed:`, error);
|
||||
});
|
||||
this.#errors.delete(botId);
|
||||
this.#touch();
|
||||
});
|
||||
return this.status();
|
||||
}
|
||||
|
||||
status() {
|
||||
const bots = this.#configStore.list().map((config) => {
|
||||
const runtimeStatus = this.#runtimes.get(config.botId)?.status ?? null;
|
||||
const connected = runtimeStatus?.ready === true
|
||||
&& runtimeStatus.wecomConnectionState === 'connected'
|
||||
&& runtimeStatus.harnessReachable === true;
|
||||
const state = connected ? 'connected'
|
||||
: runtimeStatus?.wecomConnectionState === 'connecting' ? 'connecting'
|
||||
: this.#errors.has(config.botId) || runtimeStatus?.wecomConnectionState === 'failed'
|
||||
? 'error' : 'offline';
|
||||
return {
|
||||
botId: config.botId,
|
||||
state,
|
||||
connected,
|
||||
configured: true,
|
||||
bot: { name: '企业微信机器人', appIdMasked: maskWecomBotId(config.remoteBotId) },
|
||||
health: {
|
||||
status: connected ? 'healthy' : state === 'error' ? 'error' : 'offline',
|
||||
summary: connected ? '企业微信 WebSocket 长连接运行正常'
|
||||
: state === 'error' ? '企业微信连接未就绪,插件会自动重试' : '企业微信连接当前离线',
|
||||
lastCheckedAt: runtimeStatus?.lastCheckedAt ?? null,
|
||||
lastConnectedAt: runtimeStatus?.lastConnectedAt ?? null,
|
||||
},
|
||||
stats: {
|
||||
messagesReceived: runtimeStatus?.messagesReceived ?? 0,
|
||||
messagesReplied: runtimeStatus?.messagesReplied ?? 0,
|
||||
},
|
||||
error: structuredClone(this.#errors.get(config.botId) ?? null),
|
||||
};
|
||||
});
|
||||
const connectedCount = bots.filter((bot) => bot.connected).length;
|
||||
const active = this.#activeAttemptId ? this.#attempts.get(this.#activeAttemptId) : null;
|
||||
return {
|
||||
schemaVersion: 1,
|
||||
revision: this.#revision,
|
||||
state: active && ACTIVE_ATTEMPT_STATES.has(active.state) ? 'provisioning'
|
||||
: bots.length === 0 ? 'disconnected'
|
||||
: connectedCount === bots.length ? 'connected'
|
||||
: connectedCount > 0 ? 'degraded' : 'offline',
|
||||
bots,
|
||||
totals: { configured: bots.length, connected: connectedCount },
|
||||
...(active && ACTIVE_ATTEMPT_STATES.has(active.state)
|
||||
? { provisioning: publicAttempt(active) } : {}),
|
||||
};
|
||||
}
|
||||
|
||||
async close() {
|
||||
if (this.#closed) return;
|
||||
this.#closed = true;
|
||||
if (this.#activeAttemptId) await this.cancelProvisioning(this.#activeAttemptId);
|
||||
await Promise.allSettled([...this.#runtimes.keys()].map((botId) => this.#stopRuntime(botId)));
|
||||
}
|
||||
|
||||
async #pollAttempt(record) {
|
||||
try {
|
||||
const result = await this.#qrAuth.poll({ scode: record.scode, signal: record.controller.signal });
|
||||
if (record.controller.signal.aborted || TERMINAL_ATTEMPT_STATES.has(record.state)) return;
|
||||
if (result.status === 'waiting') return;
|
||||
if (result.status === 'expired') {
|
||||
record.state = 'expired';
|
||||
record.error = safeError('expired', '企业微信二维码已过期,请重新生成。');
|
||||
this.#finishAttempt(record);
|
||||
return;
|
||||
}
|
||||
if (result.status !== 'success') {
|
||||
record.state = 'failed';
|
||||
record.error = safeError('qr-connect-failed', '企业微信扫码没有完成,请重新生成二维码。');
|
||||
this.#finishAttempt(record);
|
||||
return;
|
||||
}
|
||||
record.state = 'connecting';
|
||||
record.verificationUrl = null;
|
||||
record.expiresAt = null;
|
||||
record.scode = null;
|
||||
this.#touch();
|
||||
const transition = this.#completeProvisioning(record, result);
|
||||
record.transition = transition;
|
||||
await transition;
|
||||
} catch (error) {
|
||||
if (record.controller.signal.aborted) return;
|
||||
record.state = 'failed';
|
||||
record.error = safeError('qr-connect-failed', '企业微信扫码服务暂时不可用,请重新生成二维码。');
|
||||
this.#logger.warn?.('[dsh-im:wecom] QR polling failed');
|
||||
this.#finishAttempt(record);
|
||||
}
|
||||
}
|
||||
|
||||
async #completeProvisioning(record, result) {
|
||||
try {
|
||||
const remoteBotId = cleanString(result.remoteBotId);
|
||||
const secret = cleanString(result.secret);
|
||||
if (!remoteBotId || !secret) throw new Error('Enterprise WeChat authorization returned incomplete credentials');
|
||||
record.botId = await this.#activateBot(record, { remoteBotId, secret });
|
||||
record.state = 'connected';
|
||||
record.error = null;
|
||||
} catch (error) {
|
||||
if (record.controller.signal.aborted) {
|
||||
record.state = 'cancelled';
|
||||
record.error = safeError('cancelled', '扫码绑定已取消。');
|
||||
} else {
|
||||
record.state = 'failed';
|
||||
record.error = safeError('activation-failed', '企业微信已授权,但无法安全保存接入配置。');
|
||||
this.#logger.error?.('[dsh-im:wecom] provisioning failed');
|
||||
}
|
||||
} finally {
|
||||
this.#finishAttempt(record);
|
||||
}
|
||||
}
|
||||
|
||||
async #activateBot(record, { remoteBotId, secret }) {
|
||||
const identity = deriveWecomBotIdentity(remoteBotId);
|
||||
const previousConfig = this.#configStore.getByRemoteBotId(remoteBotId);
|
||||
const previousSecret = await this.#credentials.resolve(identity.secretRef).catch(() => undefined);
|
||||
const config = {
|
||||
botId: identity.botId,
|
||||
remoteBotId,
|
||||
secretRef: identity.secretRef,
|
||||
createdAt: previousConfig?.createdAt ?? new Date().toISOString(),
|
||||
connectedAt: new Date().toISOString(),
|
||||
};
|
||||
return this.#withBotTransition(identity.botId, async () => {
|
||||
await this.#credentials.set(identity.secretRef, secret);
|
||||
try {
|
||||
if (record.controller.signal.aborted) throw new DOMException('Cancelled', 'AbortError');
|
||||
await this.#configStore.save(config);
|
||||
} catch (error) {
|
||||
await this.#restoreCredential(identity.secretRef, previousSecret);
|
||||
throw error;
|
||||
}
|
||||
try {
|
||||
if (record.controller.signal.aborted) throw new DOMException('Cancelled', 'AbortError');
|
||||
await this.#startRuntime(config, secret);
|
||||
this.#errors.delete(identity.botId);
|
||||
} catch (error) {
|
||||
if (record.controller.signal.aborted) {
|
||||
await this.#stopRuntime(identity.botId);
|
||||
if (previousConfig) await this.#configStore.save(previousConfig).catch(() => undefined);
|
||||
else await this.#configStore.remove(identity.botId).catch(() => undefined);
|
||||
await this.#restoreCredential(identity.secretRef, previousSecret);
|
||||
throw error;
|
||||
}
|
||||
this.#errors.set(identity.botId, safeError('connection-failed', '企业微信机器人已绑定,消息连接暂未就绪。'));
|
||||
this.#logger.warn?.(`[dsh-im:wecom] bot ${identity.botId} activation connection failed`);
|
||||
}
|
||||
this.#touch();
|
||||
return identity.botId;
|
||||
});
|
||||
}
|
||||
|
||||
async #startRuntime(config, secret) {
|
||||
await this.#stopRuntime(config.botId);
|
||||
const runtime = await this.#createRuntime({ botId: config.botId, config, secret });
|
||||
if (!runtime || typeof runtime.start !== 'function' || typeof runtime.stop !== 'function') {
|
||||
throw new TypeError('createRuntime returned an invalid Enterprise WeChat runtime');
|
||||
}
|
||||
this.#runtimes.set(config.botId, runtime);
|
||||
try {
|
||||
await runtime.start();
|
||||
} catch (error) {
|
||||
await runtime.stop().catch(() => undefined);
|
||||
this.#runtimes.delete(config.botId);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async #stopRuntime(botId) {
|
||||
const runtime = this.#runtimes.get(botId);
|
||||
this.#runtimes.delete(botId);
|
||||
await runtime?.stop().catch((error) => {
|
||||
this.#logger.warn?.(`[dsh-im:wecom] bot ${botId} failed to stop cleanly:`, error);
|
||||
});
|
||||
}
|
||||
|
||||
async #resolveSecret(ref) {
|
||||
const result = await this.#credentials.resolve(ref).catch(() => undefined);
|
||||
return cleanString(result?.value);
|
||||
}
|
||||
|
||||
async #restoreCredential(ref, previous) {
|
||||
if (previous?.value) await this.#credentials.set(ref, previous.value).catch(() => undefined);
|
||||
else await this.#credentials.unset(ref).catch(() => undefined);
|
||||
}
|
||||
|
||||
#withBotTransition(botId, operation) {
|
||||
const previous = this.#transitions.get(botId) ?? Promise.resolve();
|
||||
const current = previous.catch(() => undefined).then(operation);
|
||||
const settled = current.finally(() => {
|
||||
if (this.#transitions.get(botId) === settled) this.#transitions.delete(botId);
|
||||
});
|
||||
this.#transitions.set(botId, settled);
|
||||
return settled;
|
||||
}
|
||||
|
||||
#finishAttempt(record) {
|
||||
record.scode = null;
|
||||
record.verificationUrl = null;
|
||||
record.expiresAt = null;
|
||||
if (this.#activeAttemptId === record.id) this.#activeAttemptId = null;
|
||||
this.#touch();
|
||||
const terminal = [...this.#attempts.values()].filter((attempt) => TERMINAL_ATTEMPT_STATES.has(attempt.state));
|
||||
while (terminal.length > 16) this.#attempts.delete(terminal.shift().id);
|
||||
}
|
||||
|
||||
#touch() {
|
||||
this.#revision += 1;
|
||||
}
|
||||
}
|
||||
209
src/channels/wecom/wecom-runtime.mjs
Normal file
209
src/channels/wecom/wecom-runtime.mjs
Normal file
|
|
@ -0,0 +1,209 @@
|
|||
import { WSAuthFailureError, WSClient, WSReconnectExhaustedError } from '@wecom/aibot-node-sdk';
|
||||
|
||||
import { createWecomBridgeStatus, WecomHarnessBridge } from './wecom-bridge.mjs';
|
||||
|
||||
function timeoutError() {
|
||||
const error = new Error('Enterprise WeChat WebSocket authentication timed out');
|
||||
error.code = 'connect-timeout';
|
||||
return error;
|
||||
}
|
||||
|
||||
export function createWecomRuntimeStatus() {
|
||||
return {
|
||||
startedAt: null,
|
||||
ready: false,
|
||||
wecomConnectionState: 'idle',
|
||||
harnessReachable: false,
|
||||
lastCheckedAt: null,
|
||||
lastConnectedAt: null,
|
||||
lastError: null,
|
||||
...createWecomBridgeStatus(),
|
||||
};
|
||||
}
|
||||
|
||||
export class WecomRuntime {
|
||||
#config;
|
||||
#secret;
|
||||
#harness;
|
||||
#state;
|
||||
#logger;
|
||||
#replyTimeoutMs;
|
||||
#connectTimeoutMs;
|
||||
#maxReconnectAttempts;
|
||||
#createClient;
|
||||
#status = createWecomRuntimeStatus();
|
||||
#client = null;
|
||||
#bridge = null;
|
||||
#starting = null;
|
||||
#startController = null;
|
||||
|
||||
constructor({
|
||||
config,
|
||||
secret,
|
||||
harness,
|
||||
state,
|
||||
logger = console,
|
||||
replyTimeoutMs = 600_000,
|
||||
connectTimeoutMs = 20_000,
|
||||
maxReconnectAttempts = 10,
|
||||
createClient = (options) => new WSClient(options),
|
||||
}) {
|
||||
if (!config || !secret || !harness || !state) {
|
||||
throw new TypeError('WecomRuntime requires config, secret, Harness, and state');
|
||||
}
|
||||
this.#config = config;
|
||||
this.#secret = secret;
|
||||
this.#harness = harness;
|
||||
this.#state = state;
|
||||
this.#logger = logger;
|
||||
this.#replyTimeoutMs = replyTimeoutMs;
|
||||
this.#connectTimeoutMs = connectTimeoutMs;
|
||||
this.#maxReconnectAttempts = maxReconnectAttempts;
|
||||
this.#createClient = createClient;
|
||||
}
|
||||
|
||||
get status() {
|
||||
return structuredClone(this.#status);
|
||||
}
|
||||
|
||||
async start() {
|
||||
if (this.#status.ready && this.#client) return this.status;
|
||||
if (this.#starting) return this.#starting;
|
||||
const controller = new AbortController();
|
||||
this.#startController = controller;
|
||||
this.#starting = this.#start(controller.signal).finally(() => {
|
||||
if (this.#startController === controller) this.#startController = null;
|
||||
this.#starting = null;
|
||||
});
|
||||
return this.#starting;
|
||||
}
|
||||
|
||||
async #start(signal) {
|
||||
await this.#stopActive();
|
||||
signal.throwIfAborted();
|
||||
this.#status.startedAt = new Date().toISOString();
|
||||
this.#status.wecomConnectionState = 'connecting';
|
||||
this.#status.lastError = null;
|
||||
await this.#harness.ensureRunning();
|
||||
this.#status.harnessReachable = true;
|
||||
|
||||
const silentSdkLogger = { debug() {}, info() {}, warn() {}, error() {} };
|
||||
const client = this.#createClient({
|
||||
botId: this.#config.remoteBotId,
|
||||
secret: this.#secret,
|
||||
logger: silentSdkLogger,
|
||||
maxReconnectAttempts: this.#maxReconnectAttempts,
|
||||
});
|
||||
if (!client || typeof client.connect !== 'function' || typeof client.disconnect !== 'function') {
|
||||
throw new TypeError('Enterprise WeChat client factory returned an invalid client');
|
||||
}
|
||||
this.#client = client;
|
||||
this.#bridge = new WecomHarnessBridge({
|
||||
client,
|
||||
harness: this.#harness,
|
||||
state: this.#state,
|
||||
status: this.#status,
|
||||
logger: this.#logger,
|
||||
replyTimeoutMs: this.#replyTimeoutMs,
|
||||
});
|
||||
|
||||
let readyResolve;
|
||||
let readyReject;
|
||||
let authenticated = false;
|
||||
const ready = new Promise((resolve, reject) => {
|
||||
readyResolve = resolve;
|
||||
readyReject = reject;
|
||||
});
|
||||
const onAuthenticated = () => {
|
||||
if (this.#client !== client || signal.aborted) return;
|
||||
authenticated = true;
|
||||
const now = Date.now();
|
||||
this.#status.ready = true;
|
||||
this.#status.wecomConnectionState = 'connected';
|
||||
this.#status.lastCheckedAt = now;
|
||||
this.#status.lastConnectedAt = now;
|
||||
this.#status.lastError = null;
|
||||
readyResolve();
|
||||
};
|
||||
const onDisconnected = () => {
|
||||
if (this.#client !== client) return;
|
||||
this.#status.ready = false;
|
||||
this.#status.wecomConnectionState = 'connecting';
|
||||
this.#status.lastCheckedAt = Date.now();
|
||||
};
|
||||
const onReconnecting = () => {
|
||||
if (this.#client !== client) return;
|
||||
this.#status.ready = false;
|
||||
this.#status.wecomConnectionState = 'connecting';
|
||||
this.#status.lastCheckedAt = Date.now();
|
||||
};
|
||||
const onError = (error) => {
|
||||
if (this.#client !== client) return;
|
||||
const terminal = error instanceof WSAuthFailureError || error instanceof WSReconnectExhaustedError;
|
||||
if (!authenticated && terminal) readyReject(error);
|
||||
if (terminal) {
|
||||
this.#status.ready = false;
|
||||
this.#status.wecomConnectionState = 'failed';
|
||||
}
|
||||
this.#status.lastError = terminal ? error.name : 'connection-error';
|
||||
this.#logger.warn?.(`[dsh-im:wecom] bot ${this.#config.botId} connection error`);
|
||||
};
|
||||
const onMessage = (frame) => this.#bridge?.accept(frame);
|
||||
client.on('authenticated', onAuthenticated);
|
||||
client.on('disconnected', onDisconnected);
|
||||
client.on('reconnecting', onReconnecting);
|
||||
client.on('error', onError);
|
||||
client.on('message', onMessage);
|
||||
|
||||
let timer;
|
||||
try {
|
||||
client.connect();
|
||||
await Promise.race([
|
||||
ready,
|
||||
new Promise((_, reject) => {
|
||||
signal.addEventListener('abort', () => reject(signal.reason), { once: true });
|
||||
}),
|
||||
new Promise((_, reject) => {
|
||||
timer = setTimeout(() => reject(timeoutError()), this.#connectTimeoutMs);
|
||||
}),
|
||||
]);
|
||||
return this.status;
|
||||
} catch (error) {
|
||||
if (signal.aborted) {
|
||||
await this.#stopActive();
|
||||
throw error;
|
||||
}
|
||||
this.#status.ready = false;
|
||||
this.#status.wecomConnectionState = 'failed';
|
||||
this.#status.lastError = error?.message ?? String(error);
|
||||
await this.#stopActive();
|
||||
throw error;
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
|
||||
async stop() {
|
||||
const starting = this.#starting;
|
||||
this.#startController?.abort(new DOMException('Enterprise WeChat runtime stopped', 'AbortError'));
|
||||
await this.#stopActive();
|
||||
await starting?.catch(() => undefined);
|
||||
return this.status;
|
||||
}
|
||||
|
||||
async #stopActive() {
|
||||
const client = this.#client;
|
||||
const bridge = this.#bridge;
|
||||
this.#client = null;
|
||||
this.#bridge = null;
|
||||
try {
|
||||
client?.disconnect();
|
||||
client?.removeAllListeners?.();
|
||||
} catch (error) {
|
||||
this.#logger.warn?.(`[dsh-im:wecom] bot ${this.#config.botId} failed to stop cleanly`);
|
||||
}
|
||||
await bridge?.waitForIdle();
|
||||
this.#status.ready = false;
|
||||
this.#status.wecomConnectionState = 'idle';
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue