Add QQ bot channel

This commit is contained in:
xmanrui 2026-08-15 17:51:48 +08:00
parent a422084eb5
commit 601271a7f8
33 changed files with 4971 additions and 642 deletions

View file

@ -0,0 +1,124 @@
const DEFAULT_RETRY_DELAYS_MS = Object.freeze([250, 1_000, 3_000, 5_000, 10_000, 30_000]);
function retryDelays(value) {
if (!Array.isArray(value) || value.length === 0) return [...DEFAULT_RETRY_DELAYS_MS];
const valid = value.filter((delay) => Number.isFinite(delay) && delay >= 0);
return valid.length > 0 ? valid : [...DEFAULT_RETRY_DELAYS_MS];
}
export class ConnectionSupervisor {
#controller;
#harness;
#logger;
#retryDelays;
#healthyIntervalMs;
#setTimeout;
#clearTimeout;
#timer = null;
#running = null;
#retryIndex = 0;
#closed = false;
#started = false;
#ready;
#resolveReady;
constructor({
controller,
harness,
logger = console,
retryDelaysMs,
healthyIntervalMs = 15_000,
setTimeoutImpl = setTimeout,
clearTimeoutImpl = clearTimeout,
}) {
if (!controller || typeof controller.initialize !== 'function' || typeof controller.status !== 'function') {
throw new TypeError('ConnectionSupervisor requires a controller');
}
if (!harness || typeof harness.ensureRunning !== 'function') {
throw new TypeError('ConnectionSupervisor requires a Harness client');
}
this.#controller = controller;
this.#harness = harness;
this.#logger = logger;
this.#retryDelays = retryDelays(retryDelaysMs);
this.#healthyIntervalMs = Number.isFinite(healthyIntervalMs) && healthyIntervalMs >= 0
? healthyIntervalMs : 15_000;
this.#setTimeout = setTimeoutImpl;
this.#clearTimeout = clearTimeoutImpl;
this.#ready = new Promise((resolve) => {
this.#resolveReady = resolve;
});
}
get ready() {
return this.#ready;
}
start() {
if (this.#started || this.#closed) return this;
this.#started = true;
this.#schedule(0);
return this;
}
async close() {
if (this.#closed) return;
this.#closed = true;
if (this.#timer !== null) this.#clearTimeout(this.#timer);
this.#timer = null;
await this.#running?.catch(() => undefined);
this.#resolveReady?.(null);
this.#resolveReady = null;
}
#schedule(delayMs) {
if (this.#closed) return;
this.#timer = this.#setTimeout(() => {
this.#timer = null;
void this.#run();
}, delayMs);
this.#timer?.unref?.();
}
async #run() {
if (this.#closed || this.#running) return;
const operation = this.#reconcile();
this.#running = operation;
try {
await operation;
} finally {
if (this.#running === operation) this.#running = null;
}
}
async #reconcile() {
try {
await this.#harness.ensureRunning();
if (this.#closed) return;
const status = await this.#controller.initialize();
if (this.#closed) return;
this.#resolveReady?.(status);
this.#resolveReady = null;
const { configured, connected } = status.totals;
if (connected < configured) {
const delay = this.#retryDelays[Math.min(this.#retryIndex, this.#retryDelays.length - 1)];
this.#retryIndex += 1;
this.#logger.warn?.(`[dsh-im:qq] ${connected}/${configured} bots connected; retrying in ${delay}ms`);
this.#schedule(delay);
return;
}
this.#retryIndex = 0;
this.#schedule(this.#healthyIntervalMs);
} catch (error) {
if (this.#closed) return;
const delay = this.#retryDelays[Math.min(this.#retryIndex, this.#retryDelays.length - 1)];
this.#retryIndex += 1;
this.#logger.warn?.(`[dsh-im:qq] connection reconciliation failed; retrying in ${delay}ms`, error);
this.#schedule(delay);
}
}
}
export function createConnectionSupervisor(options) {
return new ConnectionSupervisor(options);
}

View file

@ -0,0 +1,23 @@
import { createProductionController } from './production.mjs';
import { installQqRpc } from './rpc.mjs';
export const name = 'dsh-im-qq-host';
export const inject = ['connection', 'credentials', 'webServer'];
export async function apply(ctx, config = {}) {
if (config?.controller) return installQqRpc(ctx, config.controller, config.rpcOptions);
const production = await createProductionController(ctx, config, config.internals);
const disposeRpc = installQqRpc(ctx, production.controller, config.rpcOptions);
ctx.effect(() => async () => production.close(), 'dsh-im: close QQ bot connections');
return disposeRpc;
}
export function createQqHostPlugin(config) {
return Object.freeze({ name, inject, apply: (ctx) => apply(ctx, config) });
}
export { createConnectionSupervisor, ConnectionSupervisor } from './connection-supervisor.mjs';
export { createProductionController } from './production.mjs';
export { QQ_ENDPOINTS, QQ_RPC_CHANNEL, QQ_RPC_ENDPOINTS, createQqRpcHandler, installQqRpc } from './rpc.mjs';
export { QqController } from '../../../../src/channels/qq/qq-controller.mjs';
export { QqRuntime } from '../../../../src/channels/qq/qq-runtime.mjs';

View file

@ -0,0 +1,109 @@
import { unlink } from 'node:fs/promises';
import { homedir } from 'node:os';
import { join, resolve } from 'node:path';
import { QqConfigStore } from '../../../../src/channels/qq/config-store.mjs';
import { QqHarnessClient } from '../../../../src/channels/qq/harness-client.mjs';
import { QqController } from '../../../../src/channels/qq/qq-controller.mjs';
import { QqRuntime } from '../../../../src/channels/qq/qq-runtime.mjs';
import { QqQrAuth } from '../../../../src/channels/qq/qr-auth.mjs';
import { QqStateStore } from '../../../../src/channels/qq/state-store.mjs';
import { createConnectionSupervisor } from './connection-supervisor.mjs';
function harnessOrigin(webServer, configured) {
if (configured !== undefined) return new URL(configured);
const port = webServer?.port;
if (!Number.isInteger(port) || port < 1 || port > 65_535) {
throw new Error('dsh-im QQ requires an initialized DSH webServer port');
}
return new URL(`http://127.0.0.1:${port}`);
}
function pluginPaths(config) {
const dshHome = resolve(config.dshHome ?? process.env.DSH_HOME ?? join(homedir(), '.dsh'));
const root = resolve(config.dataDir ?? join(dshHome, 'integrations', 'dsh-qq'));
return {
config: resolve(config.configPath ?? join(root, 'config.json')),
bots: resolve(config.botsDir ?? join(root, 'bots')),
};
}
export async function createProductionController(ctx, config = {}, internals = {}) {
if (!ctx?.credentials) throw new TypeError('dsh-im QQ requires ctx.credentials');
if (!ctx?.webServer) throw new TypeError('dsh-im QQ requires ctx.webServer');
const ConfigStore = internals.ConfigStore ?? QqConfigStore;
const StateStore = internals.StateStore ?? QqStateStore;
const Harness = internals.HarnessClient ?? QqHarnessClient;
const Controller = internals.Controller ?? QqController;
const Runtime = internals.Runtime ?? QqRuntime;
const QrAuth = internals.QrAuth ?? QqQrAuth;
const createSupervisor = internals.createConnectionSupervisor ?? createConnectionSupervisor;
const logger = typeof ctx.logger === 'function' ? ctx.logger('dsh-im:qq') : (ctx.logger ?? console);
const paths = pluginPaths(config);
const configStore = await new ConfigStore(paths.config).load();
const qrAuth = internals.qrAuth ?? new QrAuth({ source: config.qrSource ?? 'deepseek-harness' });
const stateStores = new Map();
const statePath = (botId) => resolve(paths.bots, botId, 'state.json');
const stateFor = async (botId) => {
let state = stateStores.get(botId);
if (!state) {
state = await new StateStore(statePath(botId)).load();
stateStores.set(botId, state);
}
return state;
};
const harness = new Harness({
baseUrl: harnessOrigin(ctx.webServer, config.harnessBaseUrl),
workspace: resolve(config.workspace ?? process.cwd()),
agentPreset: config.agentPreset ?? 'standard',
autostart: false,
dshBin: config.dshBin ?? 'dsh',
});
const controller = new Controller({
qrAuth,
credentials: ctx.credentials,
configStore,
logger,
createRuntime: async ({ botId, config: botConfig, appSecret }) => new Runtime({
config: botConfig,
appSecret,
harness,
state: await stateFor(botId),
replyTimeoutMs: config.replyTimeoutMs ?? 600_000,
connectTimeoutMs: config.connectTimeoutMs ?? 20_000,
logger: {
error: (...args) => logger.error?.(`[${botId}]`, ...args),
warn: (...args) => logger.warn?.(`[${botId}]`, ...args),
info: (...args) => logger.info?.(`[${botId}]`, ...args),
debug: (...args) => logger.debug?.(`[${botId}]`, ...args),
},
}),
deleteState: async ({ botId }) => {
const state = stateStores.get(botId);
stateStores.delete(botId);
if (state && typeof state.remove === 'function') return state.remove();
try {
await unlink(statePath(botId));
} catch (error) {
if (error?.code !== 'ENOENT') throw error;
}
},
});
const supervisor = createSupervisor({
controller,
harness,
logger,
retryDelaysMs: config.retryDelaysMs,
healthyIntervalMs: config.healthyIntervalMs,
}).start();
return {
controller,
ready: supervisor.ready,
async close() {
await supervisor.close();
await controller.close();
harness.stopManagedProcess();
},
};
}

View file

@ -0,0 +1,128 @@
import QRCode from 'qrcode';
export const QQ_RPC_CHANNEL = '/qq';
export const QQ_ENDPOINTS = Object.freeze({
status: 'connection.status',
beginProvisioning: 'provision.begin',
pollProvisioning: 'provision.poll',
cancelProvisioning: 'provision.cancel',
reconnectBot: 'bot.reconnect',
deleteBot: 'bot.delete',
});
export const QQ_RPC_ENDPOINTS = Object.freeze(Object.values(QQ_ENDPOINTS));
const FORBIDDEN_PUBLIC_KEYS = new Set([
'appSecret', 'app_secret', 'secretRef', 'ownerUserOpenid', 'userOpenid', 'verificationUrl',
]);
function isRecord(value) {
return value !== null && typeof value === 'object' && !Array.isArray(value);
}
function exactKeys(value, allowed) {
return isRecord(value) && Object.keys(value).every((key) => allowed.includes(key));
}
function validId(value) {
return typeof value === 'string' && /^[A-Za-z0-9_-]{1,128}$/.test(value);
}
function payloadFailure(endpoint, payload) {
if (!isRecord(payload)) return 'Payload must be an object.';
if (endpoint === QQ_ENDPOINTS.status) return exactKeys(payload, []) ? null : 'connection.status does not accept fields.';
if (endpoint === QQ_ENDPOINTS.beginProvisioning) {
return exactKeys(payload, ['locale']) && (payload.locale === undefined || payload.locale === 'zh-CN')
? null : 'provision.begin received unsupported fields.';
}
if ([QQ_ENDPOINTS.pollProvisioning, QQ_ENDPOINTS.cancelProvisioning].includes(endpoint)) {
return exactKeys(payload, ['attemptId']) && validId(payload.attemptId)
? null : `${endpoint} requires an attemptId.`;
}
if (endpoint === QQ_ENDPOINTS.reconnectBot) {
return exactKeys(payload, ['botId']) && validId(payload.botId) ? null : 'bot.reconnect requires a botId.';
}
if (endpoint === QQ_ENDPOINTS.deleteBot) {
return exactKeys(payload, ['botId', 'confirm']) && validId(payload.botId) && payload.confirm === true
? null : 'bot.delete requires a botId and confirm=true.';
}
return 'Unknown QQ endpoint.';
}
function sanitizePublic(value) {
if (Array.isArray(value)) return value.map(sanitizePublic);
if (!isRecord(value)) return value;
const safe = {};
for (const [key, child] of Object.entries(value)) {
if (!FORBIDDEN_PUBLIC_KEYS.has(key)) safe[key] = sanitizePublic(child);
}
return safe;
}
async function qrDataUrl(value) {
return QRCode.toDataURL(value, {
type: 'image/png', errorCorrectionLevel: 'M', margin: 2, width: 320,
});
}
async function withEncodedQr(value, encodeQr) {
if (!value || typeof value.verificationUrl !== 'string') return sanitizePublic(value);
return sanitizePublic({ ...value, qrCodeDataUrl: await encodeQr(value.verificationUrl) });
}
async function publicStatus(status, encodeQr) {
const value = structuredClone(status);
if (value?.provisioning) value.provisioning = await withEncodedQr(value.provisioning, encodeQr);
return sanitizePublic(value);
}
export function createQqRpcHandler(controller, { encodeQr = qrDataUrl } = {}) {
for (const method of ['status', 'startProvisioning', 'registrationStatus', 'cancelProvisioning', 'reconnectBot', 'deleteBot']) {
if (typeof controller?.[method] !== 'function') throw new TypeError(`A complete QQ controller is required (${method})`);
}
const qrCache = new Map();
const cachedEncode = (url) => {
let encoded = qrCache.get(url);
if (!encoded) {
if (qrCache.size >= 16) qrCache.delete(qrCache.keys().next().value);
encoded = Promise.resolve().then(() => encodeQr(url));
qrCache.set(url, encoded);
}
return encoded;
};
return async (endpoint, payload, signal) => {
if (signal?.aborted) return { ok: false, error: { code: 'cancelled', message: 'The request was cancelled.' } };
if (!QQ_RPC_ENDPOINTS.includes(endpoint)) return { ok: false, error: { code: 'bad-request', message: 'Unknown QQ endpoint.' } };
const invalid = payloadFailure(endpoint, payload);
if (invalid) return { ok: false, error: { code: 'bad-request', message: invalid } };
try {
let value;
if (endpoint === QQ_ENDPOINTS.status) value = await publicStatus(await controller.status(), cachedEncode);
else if (endpoint === QQ_ENDPOINTS.beginProvisioning) value = await withEncodedQr(await controller.startProvisioning(), cachedEncode);
else if (endpoint === QQ_ENDPOINTS.pollProvisioning) {
const current = await controller.registrationStatus(payload.attemptId);
if (!current) return { ok: false, error: { code: 'bad-request', message: 'The provisioning attempt no longer exists.' } };
value = await withEncodedQr(current, cachedEncode);
} else if (endpoint === QQ_ENDPOINTS.cancelProvisioning) {
value = sanitizePublic(await controller.cancelProvisioning(payload.attemptId));
} else if (endpoint === QQ_ENDPOINTS.reconnectBot) {
value = await publicStatus(await controller.reconnectBot(payload.botId), cachedEncode);
} else {
value = await publicStatus(await controller.deleteBot(payload.botId), cachedEncode);
}
return signal?.aborted
? { ok: false, error: { code: 'cancelled', message: 'The request was cancelled.' } }
: { ok: true, value };
} catch {
return signal?.aborted
? { ok: false, error: { code: 'cancelled', message: 'The request was cancelled.' } }
: { ok: false, error: { code: 'qq-operation-failed', message: 'QQ 操作失败,请稍后重试。' } };
}
};
}
export function installQqRpc(ctx, controller, options) {
if (!ctx?.connection?.rpc || typeof ctx.connection.rpc.handle !== 'function') {
throw new TypeError('DSH Host Connection RPC is required');
}
return ctx.connection.rpc.handle(QQ_RPC_CHANNEL, createQqRpcHandler(controller, options), { authority: 'loopback' });
}

View file

@ -1,5 +1,6 @@
import { apply as applyDingtalk } from './channels/dingtalk/index.mjs';
import { apply as applyFeishu } from './channels/feishu/index.mjs';
import { apply as applyQq } from './channels/qq/index.mjs';
import { apply as applyWeixin } from './channels/weixin/index.mjs';
export const name = 'dsh-im-host';
@ -9,6 +10,7 @@ export function createImHostPlugin(internals = {}) {
const startFeishu = internals.applyFeishu ?? applyFeishu;
const startWeixin = internals.applyWeixin ?? applyWeixin;
const startDingtalk = internals.applyDingtalk ?? applyDingtalk;
const startQq = internals.applyQq ?? applyQq;
return Object.freeze({
name,
inject,
@ -16,6 +18,7 @@ export function createImHostPlugin(internals = {}) {
await startFeishu(ctx, config.feishu ?? {});
await startWeixin(ctx, config.weixin ?? {});
await startDingtalk(ctx, config.dingtalk ?? {});
await startQq(ctx, config.qq ?? {});
},
});
}