Internalize all IM channel implementations

This commit is contained in:
xmanrui 2026-08-15 15:40:53 +08:00
parent 8996476693
commit bd469f58f8
95 changed files with 25012 additions and 100 deletions

View file

@ -0,0 +1,130 @@
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 delayMs = this.#retryDelays[Math.min(this.#retryIndex, this.#retryDelays.length - 1)];
this.#retryIndex += 1;
this.#logger.warn?.(
`[dsh-dingtalk] ${connected}/${configured} bots connected; retrying in ${delayMs}ms`,
);
this.#schedule(delayMs);
return;
}
this.#retryIndex = 0;
this.#schedule(this.#healthyIntervalMs);
} catch (error) {
if (this.#closed) return;
const delayMs = this.#retryDelays[Math.min(this.#retryIndex, this.#retryDelays.length - 1)];
this.#retryIndex += 1;
this.#logger.warn?.(
`[dsh-dingtalk] connection reconciliation failed; retrying in ${delayMs}ms`,
error,
);
this.#schedule(delayMs);
}
}
}
export function createConnectionSupervisor(options) {
return new ConnectionSupervisor(options);
}

View file

@ -0,0 +1,32 @@
import { createProductionController } from './production.mjs';
import { installDingtalkRpc } from './rpc.mjs';
export const name = 'dsh-dingtalk-host';
export const inject = ['connection', 'credentials', 'webServer'];
export async function apply(ctx, config = {}) {
if (config?.controller) return installDingtalkRpc(ctx, config.controller, config.rpcOptions);
const production = await createProductionController(ctx, config, config.internals);
const disposeRpc = installDingtalkRpc(ctx, production.controller, config.rpcOptions);
ctx.effect(() => async () => {
await production.close();
}, 'dsh-dingtalk: close bot connections');
return disposeRpc;
}
export function createDingtalkHostPlugin(config) {
return Object.freeze({ name, inject, apply: (ctx) => apply(ctx, config) });
}
export { createConnectionSupervisor, ConnectionSupervisor } from './connection-supervisor.mjs';
export { createProductionController } from './production.mjs';
export {
DINGTALK_ENDPOINTS,
DINGTALK_RPC_CHANNEL,
DINGTALK_RPC_ENDPOINTS,
createDingtalkRpcHandler,
installDingtalkRpc,
} from './rpc.mjs';
export { DingtalkController } from '../../../../src/channels/dingtalk/dingtalk-controller.mjs';
export { DingtalkRuntime } from '../../../../src/channels/dingtalk/dingtalk-runtime.mjs';

View file

@ -0,0 +1,122 @@
import { unlink } from 'node:fs/promises';
import { homedir } from 'node:os';
import { join, resolve } from 'node:path';
import { DingtalkConfigStore } from '../../../../src/channels/dingtalk/config-store.mjs';
import { DingtalkDeviceAuth } from '../../../../src/channels/dingtalk/device-auth.mjs';
import { DingtalkController } from '../../../../src/channels/dingtalk/dingtalk-controller.mjs';
import { DingtalkRuntime } from '../../../../src/channels/dingtalk/dingtalk-runtime.mjs';
import { HarnessClient } from '../../../../src/channels/dingtalk/harness-client.mjs';
import { DingtalkStateStore } from '../../../../src/channels/dingtalk/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-dingtalk 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-dingtalk'));
return {
root,
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-dingtalk requires ctx.credentials');
if (!ctx?.webServer) throw new TypeError('dsh-dingtalk requires ctx.webServer');
const ConfigStore = internals.ConfigStore ?? DingtalkConfigStore;
const DeviceAuth = internals.DeviceAuth ?? DingtalkDeviceAuth;
const StateStore = internals.StateStore ?? DingtalkStateStore;
const Harness = internals.HarnessClient ?? HarnessClient;
const Controller = internals.Controller ?? DingtalkController;
const Runtime = internals.Runtime ?? DingtalkRuntime;
const createSupervisor = internals.createConnectionSupervisor ?? createConnectionSupervisor;
const logger = typeof ctx.logger === 'function'
? ctx.logger('dsh-dingtalk')
: (ctx.logger ?? console);
const paths = pluginPaths(config);
const configStore = await new ConfigStore(paths.config).load();
const deviceAuth = internals.deviceAuth ?? new DeviceAuth({
baseUrl: config.registrationBaseUrl,
});
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({
deviceAuth,
credentials: ctx.credentials,
configStore,
logger,
createRuntime: async ({ botId, config: botConfig, clientSecret }) => {
const state = await stateFor(botId);
return new Runtime({
config: botConfig,
clientSecret,
harness,
state,
replyTimeoutMs: config.replyTimeoutMs ?? 600_000,
maxMessageChars: config.maxMessageChars ?? 4_000,
connectTimeoutMs: config.connectTimeoutMs ?? 15_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') {
await state.remove();
return;
}
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,221 @@
import QRCode from 'qrcode';
export const DINGTALK_RPC_CHANNEL = '/dingtalk';
export const DINGTALK_ENDPOINTS = Object.freeze({
status: 'connection.status',
beginProvisioning: 'provision.begin',
pollProvisioning: 'provision.poll',
cancelProvisioning: 'provision.cancel',
reconnectBot: 'bot.reconnect',
deleteBot: 'bot.delete',
approveSender: 'bot.sender.approve',
revokeSender: 'bot.sender.revoke',
});
export const DINGTALK_RPC_ENDPOINTS = Object.freeze(Object.values(DINGTALK_ENDPOINTS));
const FORBIDDEN_PUBLIC_KEYS = new Set([
'clientSecret',
'client_secret',
'deviceCode',
'device_code',
'secretRef',
'staffId',
'senderStaffId',
'verificationUrl',
'verificationUri',
'userCode',
]);
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 === DINGTALK_ENDPOINTS.status) {
return exactKeys(payload, []) ? null : 'connection.status does not accept fields.';
}
if (endpoint === DINGTALK_ENDPOINTS.beginProvisioning) {
return exactKeys(payload, ['locale']) && (payload.locale === undefined || payload.locale === 'zh-CN')
? null
: 'provision.begin received unsupported fields.';
}
if ([DINGTALK_ENDPOINTS.pollProvisioning, DINGTALK_ENDPOINTS.cancelProvisioning].includes(endpoint)) {
return exactKeys(payload, ['attemptId']) && validId(payload.attemptId)
? null
: `${endpoint} requires an attemptId.`;
}
if (endpoint === DINGTALK_ENDPOINTS.reconnectBot) {
return exactKeys(payload, ['botId']) && validId(payload.botId)
? null
: 'bot.reconnect requires a botId.';
}
if (endpoint === DINGTALK_ENDPOINTS.deleteBot) {
return exactKeys(payload, ['botId', 'confirm']) && validId(payload.botId) && payload.confirm === true
? null
: 'bot.delete requires a botId and confirm=true.';
}
if (endpoint === DINGTALK_ENDPOINTS.approveSender) {
return exactKeys(payload, ['botId', 'requestId', 'confirm'])
&& validId(payload.botId)
&& validId(payload.requestId)
&& payload.confirm === true
? null
: 'bot.sender.approve requires botId, requestId, and confirm=true.';
}
if (endpoint === DINGTALK_ENDPOINTS.revokeSender) {
return exactKeys(payload, ['botId', 'senderKey', 'confirm'])
&& validId(payload.botId)
&& validId(payload.senderKey)
&& payload.confirm === true
? null
: 'bot.sender.revoke requires botId, senderKey, and confirm=true.';
}
return 'Unknown DingTalk endpoint.';
}
function badRequest(message) {
return { ok: false, error: { code: 'bad-request', message } };
}
function cancelled() {
return { ok: false, error: { code: 'cancelled', message: 'The request was cancelled.' } };
}
function internalFailure() {
return {
ok: false,
error: { code: 'dingtalk-operation-failed', message: '钉钉操作失败,请稍后重试。' },
};
}
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);
}
function assertController(controller) {
for (const method of [
'status',
'startProvisioning',
'registrationStatus',
'cancelProvisioning',
'reconnectBot',
'deleteBot',
'approveSender',
'revokeSender',
]) {
if (typeof controller?.[method] !== 'function') {
throw new TypeError(`A complete DingTalk controller is required (${method})`);
}
}
}
export function createDingtalkRpcHandler(controller, { encodeQr = qrDataUrl } = {}) {
assertController(controller);
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 cancelled();
if (!DINGTALK_RPC_ENDPOINTS.includes(endpoint)) return badRequest('Unknown DingTalk endpoint.');
const invalid = payloadFailure(endpoint, payload);
if (invalid) return badRequest(invalid);
try {
let value;
if (endpoint === DINGTALK_ENDPOINTS.status) {
value = await publicStatus(await controller.status(), cachedEncode);
} else if (endpoint === DINGTALK_ENDPOINTS.beginProvisioning) {
const started = await controller.startProvisioning({ signal });
if (signal?.aborted) {
await controller.cancelProvisioning(started.attemptId);
return cancelled();
}
value = await withEncodedQr(started, cachedEncode);
} else if (endpoint === DINGTALK_ENDPOINTS.pollProvisioning) {
const current = await controller.registrationStatus(payload.attemptId);
if (!current) return badRequest('The provisioning attempt no longer exists.');
value = await withEncodedQr(current, cachedEncode);
} else if (endpoint === DINGTALK_ENDPOINTS.cancelProvisioning) {
value = await controller.cancelProvisioning(payload.attemptId);
if (!value) return badRequest('The provisioning attempt no longer exists.');
value = sanitizePublic(value);
} else if (endpoint === DINGTALK_ENDPOINTS.reconnectBot) {
value = await publicStatus(await controller.reconnectBot(payload.botId), cachedEncode);
} else if (endpoint === DINGTALK_ENDPOINTS.deleteBot) {
value = await publicStatus(await controller.deleteBot(payload.botId), cachedEncode);
} else if (endpoint === DINGTALK_ENDPOINTS.approveSender) {
value = await publicStatus(
await controller.approveSender(payload.botId, payload.requestId),
cachedEncode,
);
} else {
value = await publicStatus(
await controller.revokeSender(payload.botId, payload.senderKey),
cachedEncode,
);
}
return signal?.aborted ? cancelled() : { ok: true, value };
} catch {
return signal?.aborted ? cancelled() : internalFailure();
}
};
}
export function installDingtalkRpc(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(
DINGTALK_RPC_CHANNEL,
createDingtalkRpcHandler(controller, options),
{ authority: 'loopback' },
);
}

View file

@ -0,0 +1,167 @@
const DEFAULT_RETRY_DELAYS_MS = Object.freeze([250, 1_000, 3_000, 5_000, 10_000, 30_000]);
function safeDelay(value, fallback) {
return Number.isFinite(value) && value >= 0 ? value : fallback;
}
function safeRetryDelays(value) {
if (!Array.isArray(value) || value.length === 0) return [...DEFAULT_RETRY_DELAYS_MS];
const delays = value.map((delay) => safeDelay(delay, -1)).filter((delay) => delay >= 0);
return delays.length > 0 ? delays : [...DEFAULT_RETRY_DELAYS_MS];
}
function totals(status) {
const configured = Number.isInteger(status?.totals?.configured)
? status.totals.configured
: (Array.isArray(status?.bots) ? status.bots.length : 0);
const connected = Number.isInteger(status?.totals?.connected)
? status.totals.connected
: (Array.isArray(status?.bots)
? status.bots.filter((bot) => bot?.connected === true).length
: 0);
return { configured, connected };
}
/**
* Starts bot connections only after the in-process Harness HTTP API is ready,
* then periodically reconciles failed/offline bots. Timers never keep the Host
* alive on their own and shutdown waits for an in-flight reconciliation.
*/
export class ConnectionSupervisor {
#controller;
#harness;
#logger;
#retryDelaysMs;
#healthyIntervalMs;
#setTimeout;
#clearTimeout;
#timer = null;
#running = null;
#retryIndex = 0;
#started = false;
#closed = 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.#retryDelaysMs = safeRetryDelays(retryDelaysMs);
this.#healthyIntervalMs = safeDelay(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();
} catch (error) {
if (this.#closed) return;
this.#retry('Harness Host is not ready', error);
return;
}
if (this.#closed) return;
try {
await this.#controller.initialize();
if (this.#closed) return;
const status = this.#controller.status();
this.#resolveReady?.(status);
this.#resolveReady = null;
const current = totals(status);
if (current.connected < current.configured) {
const delay = this.#retryDelaysMs[Math.min(this.#retryIndex, this.#retryDelaysMs.length - 1)];
this.#retryIndex += 1;
this.#logger.warn?.(
`[dsh-feishu] ${current.connected}/${current.configured} bots connected; retrying automatically in ${delay}ms`,
);
this.#schedule(delay);
return;
}
this.#retryIndex = 0;
this.#schedule(this.#healthyIntervalMs);
} catch (error) {
if (this.#closed) return;
this.#retry('Bot connection reconciliation failed', error);
}
}
#retry(message, error) {
const delay = this.#retryDelaysMs[Math.min(this.#retryIndex, this.#retryDelaysMs.length - 1)];
this.#retryIndex += 1;
this.#logger.warn?.(
`[dsh-feishu] ${message}; retrying automatically in ${delay}ms`,
error,
);
this.#schedule(delay);
}
}
export function createConnectionSupervisor(options) {
return new ConnectionSupervisor(options);
}

View file

@ -0,0 +1,172 @@
const ACTIVE_REGISTRATION_STATES = new Set([
'starting',
'qr_ready',
'polling',
'slow_down',
'domain_switched',
]);
function credentialResult(result) {
const appId = result?.client_id ?? result?.appId;
const appSecret = result?.client_secret ?? result?.appSecret;
if (typeof appId !== 'string' || appId.length === 0
|| typeof appSecret !== 'string' || appSecret.length === 0) {
throw new TypeError('Feishu registration returned invalid credentials');
}
return {
appId,
appSecret,
userInfo: result?.user_info ?? result?.userInfo,
};
}
async function readConnectionStatus(connectionManager) {
if (typeof connectionManager.status !== 'function') return {};
return await connectionManager.status();
}
function isConnected(status) {
if (status?.connected === true) return true;
return status?.ready === true
&& status?.feishuLongConnectionState === 'connected'
&& status?.harnessReachable === true;
}
/**
* Minimal orchestration boundary between QR provisioning and the live bot.
*
* `createProvisioningManager` receives the only callback that can observe the
* App Secret. The callback persists credentials and starts the long-lived
* connection before the provisioning manager may report `succeeded`.
*/
export class ProvisioningBackedController {
#credentialStore;
#connectionManager;
#registrationOptions;
#manager;
#knownConfigured = false;
#lastError = null;
constructor({
createProvisioningManager,
credentialStore,
connectionManager,
registrationOptions = {},
} = {}) {
if (typeof createProvisioningManager !== 'function') {
throw new TypeError('createProvisioningManager is required');
}
if (!credentialStore
|| typeof credentialStore.save !== 'function'
|| typeof credentialStore.clear !== 'function') {
throw new TypeError('credentialStore.save/clear are required');
}
if (!connectionManager
|| typeof connectionManager.connect !== 'function'
|| typeof connectionManager.disconnect !== 'function') {
throw new TypeError('connectionManager.connect/disconnect are required');
}
if (registrationOptions === null
|| typeof registrationOptions !== 'object'
|| Array.isArray(registrationOptions)) {
throw new TypeError('registrationOptions must be an object');
}
this.#credentialStore = credentialStore;
this.#connectionManager = connectionManager;
this.#registrationOptions = structuredClone(registrationOptions);
this.#manager = createProvisioningManager({
onCredentials: (result) => this.#acceptCredentials(result),
});
if (!this.#manager
|| typeof this.#manager.start !== 'function'
|| typeof this.#manager.status !== 'function'
|| typeof this.#manager.cancel !== 'function') {
throw new TypeError('The provisioning manager must implement start/status/cancel');
}
}
async startRegistration() {
this.#lastError = null;
await this.#manager.start(structuredClone(this.#registrationOptions));
return this.status();
}
async cancelRegistration() {
await this.#manager.cancel();
return this.status();
}
async disconnect() {
await this.#manager.cancel();
await this.#connectionManager.disconnect();
try {
await this.#credentialStore.clear();
this.#knownConfigured = false;
this.#lastError = null;
} catch {
// Do not reflect a credential provider's error text across the RPC
// boundary. It may include a backend path or secret reference detail.
this.#lastError = {
code: 'credential_removal_failed',
message: 'Unable to remove the Feishu credentials.',
};
}
return this.status();
}
async status() {
const registration = await this.#manager.status();
const connection = await readConnectionStatus(this.#connectionManager);
const connected = isConnected(connection);
let configured = this.#knownConfigured;
if (typeof this.#credentialStore.configured === 'function') {
try {
configured = await this.#credentialStore.configured();
} catch {
configured = this.#knownConfigured;
}
}
let phase = 'unconfigured';
if (connected) phase = 'connected';
else if (ACTIVE_REGISTRATION_STATES.has(registration?.state)) phase = 'registering';
else if (registration?.state === 'saving') phase = 'connecting';
else if (this.#lastError || registration?.state === 'error') phase = 'error';
else if (configured) phase = 'disconnected';
return {
phase,
connected,
configured,
registration,
connection,
error: this.#lastError ?? registration?.error ?? null,
};
}
async close() {
await this.#manager.cancel();
await this.#connectionManager.disconnect();
}
async #acceptCredentials(result) {
const credentials = credentialResult(result);
try {
await this.#credentialStore.save(credentials);
this.#knownConfigured = true;
await this.#connectionManager.connect(credentials);
this.#lastError = null;
} catch {
this.#lastError = {
code: 'connection_failed',
message: 'The bot was created, but its connection could not be started.',
};
throw new Error('Unable to activate the Feishu connection.');
}
}
}
export function createProvisioningBackedController(options) {
return new ProvisioningBackedController(options);
}

View file

@ -0,0 +1,84 @@
/**
* Credential references owned by the Feishu Host plugin. They deliberately
* use DSH's credential provider instead of plugin settings, so the browser's
* configuration plane can only observe configured/source metadata.
*/
export const FEISHU_APP_ID_REF = 'DSH_FEISHU_APP_ID';
export const FEISHU_APP_SECRET_REF = 'DSH_FEISHU_APP_SECRET';
function assertNonEmptyString(value, label) {
if (typeof value !== 'string' || value.length === 0) {
throw new TypeError(`${label} must be a non-empty string`);
}
return value;
}
async function restore(provider, ref, previous) {
try {
if (previous?.value) await provider.set(ref, previous.value);
else await provider.unset(ref);
} catch {
// Preserve the original write failure. The provider remains the source
// of truth and will report the partial state through describe().
}
}
/**
* Adapt the real DSH `ctx.credentials` seam to the small interface consumed
* by the Feishu controller. Secret values only travel Host-to-Host here.
*
* @param {{resolve(Function), describe(Function), set(Function), unset(Function)}} provider
* @param {{appIdRef?: string, appSecretRef?: string}} options
*/
export function createDshCredentialStore(provider, options = {}) {
if (!provider
|| typeof provider.resolve !== 'function'
|| typeof provider.describe !== 'function'
|| typeof provider.set !== 'function'
|| typeof provider.unset !== 'function') {
throw new TypeError('A DSH credential provider is required');
}
const appIdRef = options.appIdRef ?? FEISHU_APP_ID_REF;
const appSecretRef = options.appSecretRef ?? FEISHU_APP_SECRET_REF;
return Object.freeze({
async save({ appId, appSecret }) {
const nextId = assertNonEmptyString(appId, 'Feishu App ID');
const nextSecret = assertNonEmptyString(appSecret, 'Feishu App Secret');
const [previousId, previousSecret] = await Promise.all([
provider.resolve(appIdRef),
provider.resolve(appSecretRef),
]);
try {
// Store the secret first. An App ID without its matching secret must
// never be treated as a usable integration.
await provider.set(appSecretRef, nextSecret);
await provider.set(appIdRef, nextId);
} catch {
await restore(provider, appSecretRef, previousSecret);
await restore(provider, appIdRef, previousId);
throw new Error('Unable to store the Feishu credentials.');
}
},
async clear() {
const outcomes = await Promise.allSettled([
provider.unset(appIdRef),
provider.unset(appSecretRef),
]);
if (outcomes.some((outcome) => outcome.status === 'rejected')) {
throw new Error('Unable to remove the Feishu credentials.');
}
},
async configured() {
const [appId, appSecret] = await Promise.all([
provider.describe(appIdRef),
provider.describe(appSecretRef),
]);
return appId.configured === true && appSecret.configured === true;
},
});
}

View file

@ -0,0 +1,66 @@
import { createProvisioningBackedController } from './controller.mjs';
import { createProductionController } from './production.mjs';
import { installFeishuRpc } from './rpc.mjs';
export const name = 'dsh-feishu-host';
export const inject = ['connection', 'credentials', 'webServer'];
function controllerFrom(ctx, config) {
if (config?.controller) return config.controller;
if (typeof config?.createController === 'function') return config.createController();
if (typeof config?.createProvisioningManager === 'function') {
return createProvisioningBackedController(config);
}
// Cordis deliberately throws when a plugin reads an undeclared service,
// even through optional chaining. Production uses the explicit services in
// `inject`; test/programmatic controllers must therefore come from config.
return undefined;
}
/**
* Cordis/DSH Host plugin entry. Production composition may supply an owned
* controller service; tests and embedded distributions may inject one through
* config without changing the Connection RPC boundary.
*/
export async function apply(ctx, config = {}) {
const controller = controllerFrom(ctx, config);
if (controller) return installFeishuRpc(ctx, controller);
const production = await createProductionController(ctx, config);
const disposeRpc = installFeishuRpc(ctx, production.controller);
ctx.effect(() => async () => {
await production.close();
}, 'dsh-feishu: close controller and live connection');
return disposeRpc;
}
/** Create a programmatic plugin module with dependencies closed over. */
export function createFeishuHostPlugin(config) {
return Object.freeze({
name,
inject,
apply: (ctx) => apply(ctx, config),
});
}
export {
ProvisioningBackedController,
createProvisioningBackedController,
} from './controller.mjs';
export { createProductionController } from './production.mjs';
export { ConnectionSupervisor, createConnectionSupervisor } from './connection-supervisor.mjs';
export { MultiBotDshFeishuController } from '../../../../src/channels/feishu/multi-bot-controller.mjs';
export {
FEISHU_APP_ID_REF,
FEISHU_APP_SECRET_REF,
createDshCredentialStore,
} from './credential-store.mjs';
export {
FEISHU_ENDPOINTS,
FEISHU_MULTI_ENDPOINTS,
FEISHU_RPC_CHANNEL,
FEISHU_RPC_ENDPOINTS,
createFeishuRpcHandler,
installFeishuRpc,
toPublicFeishuStatus,
} from './rpc.mjs';

View file

@ -0,0 +1,138 @@
import { homedir } from 'node:os';
import { join, resolve } from 'node:path';
import { unlink } from 'node:fs/promises';
import * as Lark from '@larksuiteoapi/node-sdk';
import { createConnectionSupervisor } from './connection-supervisor.mjs';
import { verifyFeishuApp } from '../../../../src/channels/feishu/feishu-app.mjs';
import { FeishuRuntime } from '../../../../src/channels/feishu/feishu-runtime.mjs';
import { HarnessClient } from '../../../../src/channels/feishu/harness-client.mjs';
import {
LEGACY_FEISHU_SECRET_REF,
PluginConfigStore,
} from '../../../../src/channels/feishu/plugin-config-store.mjs';
import { MultiBotDshFeishuController } from '../../../../src/channels/feishu/multi-bot-controller.mjs';
import { StateStore } from '../../../../src/channels/feishu/state-store.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-feishu requires an initialized DSH webServer port');
}
// Even when DSH listens on all interfaces, its own plugin talks through the
// loopback authority accepted by Connection's request-trust fence.
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-feishu'));
return {
root,
config: resolve(config.configPath ?? join(root, 'config.json')),
legacyState: resolve(config.statePath ?? join(root, 'state.json')),
bots: resolve(config.botsDir ?? join(root, 'bots')),
};
}
/**
* Assemble the production controller from DSH Host services and local plugin
* classes. No bridge process or environment App credentials are required.
*/
export async function createProductionController(ctx, config = {}, internals = {}) {
if (!ctx?.credentials) throw new TypeError('dsh-feishu requires ctx.credentials');
if (!ctx?.webServer) throw new TypeError('dsh-feishu requires ctx.webServer');
const lark = internals.lark ?? Lark;
const Controller = internals.Controller ?? MultiBotDshFeishuController;
const ConfigStore = internals.ConfigStore ?? PluginConfigStore;
const SessionStateStore = internals.StateStore ?? StateStore;
const Harness = internals.HarnessClient ?? HarnessClient;
const Runtime = internals.FeishuRuntime ?? FeishuRuntime;
const verifyApp = internals.verifyFeishuApp ?? verifyFeishuApp;
const createSupervisor = internals.createConnectionSupervisor ?? createConnectionSupervisor;
const logger = typeof ctx.logger === 'function'
? ctx.logger('dsh-feishu')
: (ctx.logger ?? console);
const paths = pluginPaths(config);
const configStore = await new ConfigStore(paths.config).load();
// State is lazy per bot. A corrupt legacy file can therefore fail only the
// migrated bot and cannot prevent healthy v2 bots from starting.
const stateStores = new Map();
const statePathFor = (botConfig) => !botConfig.id
|| !botConfig.secretRef
|| botConfig.secretRef === LEGACY_FEISHU_SECRET_REF
? paths.legacyState
: resolve(paths.bots, botConfig.id, 'state.json');
const stateFor = async (botConfig) => {
const stateKey = botConfig.id ?? '__legacy__';
let state = stateStores.get(stateKey);
if (!state) {
state = await new SessionStateStore(statePathFor(botConfig)).load();
stateStores.set(stateKey, state);
}
return state;
};
const harness = new Harness({
baseUrl: harnessOrigin(ctx.webServer, config.harnessBaseUrl),
workspace: resolve(config.workspace ?? process.cwd()),
agentPreset: config.agentPreset ?? 'standard',
// This plugin is already hosted by a running DSH process. Starting a
// second DSH would create a competing server and lifecycle.
autostart: false,
dshBin: config.dshBin ?? 'dsh',
});
const controller = new Controller({
registerApp: (options) => lark.registerApp(options),
verifyApp,
credentials: ctx.credentials,
configStore,
createRuntime: async ({ botId, config: botConfig, appSecret }) => {
const state = await stateFor(botConfig);
return new Runtime({
lark,
appId: botConfig.appId,
appSecret,
domain: botConfig.domain,
ownerOpenIds: botConfig.ownerOpenIds ?? [botConfig.ownerOpenId],
harness,
state,
replyTimeoutMs: config.replyTimeoutMs ?? 600_000,
logger: {
error: (...args) => logger.error?.(`[${botId ?? botConfig.id}]`, ...args),
warn: (...args) => logger.warn?.(`[${botId ?? botConfig.id}]`, ...args),
info: (...args) => logger.info?.(`[${botId ?? botConfig.id}]`, ...args),
debug: (...args) => logger.debug?.(`[${botId ?? botConfig.id}]`, ...args),
},
});
},
deleteState: async ({ botId, config: botConfig }) => {
stateStores.delete(botId);
try {
await unlink(statePathFor(botConfig));
} 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,434 @@
import QRCode from 'qrcode';
import {
FEISHU_ENDPOINTS,
FEISHU_RPC_CHANNEL,
} from '../../../client/channels/feishu/api.js';
export { FEISHU_ENDPOINTS, FEISHU_RPC_CHANNEL };
export const FEISHU_MULTI_ENDPOINTS = Object.freeze({
reconnectBot: 'bot.reconnect',
disconnectBot: 'bot.disconnect',
deleteBot: 'bot.delete',
});
export const FEISHU_RPC_ENDPOINTS = Object.freeze([
...new Set([...Object.values(FEISHU_ENDPOINTS), ...Object.values(FEISHU_MULTI_ENDPOINTS)]),
]);
const REGISTRATION_STATES = new Set([
'idle', 'starting', 'qr_ready', 'polling', 'slow_down',
'domain_switched', 'saving', 'succeeded', 'expired', 'cancelled', 'error',
]);
const SAFE_ID = /^[A-Za-z0-9_-]{1,128}$/;
const PUBLIC_ERROR_MESSAGES = Object.freeze({
abort: 'Registration was cancelled.',
access_denied: 'Registration was denied.',
expired_token: 'The registration QR code expired.',
invalid_credentials: 'Feishu returned invalid app credentials.',
credentials_callback_failed: 'Unable to activate the Feishu connection.',
registration_failed: 'Unable to register the Feishu app.',
connection_failed: 'The bot was created, but its connection could not be started.',
credential_removal_failed: 'Unable to remove the Feishu credentials.',
state_cleanup_failed: 'Unable to remove the bot session data. Please retry.',
deletion_pending: 'Bot deletion is incomplete. Retry removal to finish cleanup.',
missing_credentials: 'The bot credentials are missing. Delete it and scan again.',
});
const POLL_STATUS_BY_REGISTRATION = Object.freeze({
idle: 'pending', starting: 'pending', qr_ready: 'pending', polling: 'pending',
slow_down: 'pending', domain_switched: 'pending', saving: 'connecting',
succeeded: 'connected', expired: 'expired', cancelled: 'failed', error: 'failed',
});
function isPlainObject(value) {
if (value === null || typeof value !== 'object' || Array.isArray(value)) return false;
const prototype = Object.getPrototypeOf(value);
return prototype === Object.prototype || prototype === null;
}
function hasOnlyKeys(value, allowed) {
return isPlainObject(value)
&& Reflect.ownKeys(value).every((key) => typeof key === 'string' && allowed.has(key));
}
function finiteNumber(value) {
return Number.isFinite(value) ? value : undefined;
}
function safeOpaqueId(value) {
return typeof value === 'string' && SAFE_ID.test(value);
}
function publicError(error) {
if (!error || typeof error !== 'object') return null;
const code = typeof error.code === 'string' && Object.hasOwn(PUBLIC_ERROR_MESSAGES, error.code)
? error.code
: 'registration_failed';
return { code, message: PUBLIC_ERROR_MESSAGES[code] };
}
function publicRegistration(registration) {
if (!registration || typeof registration !== 'object') return { state: 'idle', attempt: 0 };
const state = REGISTRATION_STATES.has(registration.state) ? registration.state : 'error';
const attempt = safeOpaqueId(registration.attempt)
? registration.attempt
: (finiteNumber(registration.attempt) ?? 0);
const result = { state, attempt };
const updatedAt = finiteNumber(registration.updatedAt);
const expiresAt = finiteNumber(registration.expiresAt);
const remainingSeconds = finiteNumber(registration.remainingSeconds);
const pollIntervalSeconds = finiteNumber(registration.pollIntervalSeconds);
if (updatedAt !== undefined) result.updatedAt = updatedAt;
if (typeof registration.qrCodeUrl === 'string' && registration.qrCodeUrl.length > 0) {
result.qrCodeUrl = registration.qrCodeUrl;
}
if (expiresAt !== undefined) result.expiresAt = expiresAt;
if (remainingSeconds !== undefined) result.remainingSeconds = remainingSeconds;
if (pollIntervalSeconds !== undefined) result.pollIntervalSeconds = pollIntervalSeconds;
if (safeOpaqueId(registration.botId)) result.botId = registration.botId;
const error = publicError(registration.error);
if (error) result.error = error;
return result;
}
function connectionFacts(connection) {
const source = connection && typeof connection === 'object' ? connection : {};
const connected = source.connected === true
|| (source.ready === true
&& source.feishuLongConnectionState === 'connected'
&& source.harnessReachable === true);
return {
connected,
ready: source.ready === true,
harnessReachable: source.harnessReachable === true,
};
}
function publicBot(bot) {
const source = bot && typeof bot === 'object' ? bot : {};
const result = {
name: typeof source.name === 'string' && source.name.length > 0 ? source.name : '飞书机器人',
};
if (typeof source.avatarUrl === 'string') result.avatarUrl = source.avatarUrl;
if (typeof source.appIdMasked === 'string') result.appIdMasked = source.appIdMasked;
if (typeof source.tenantName === 'string') result.tenantName = source.tenantName;
if (source.domain === 'feishu' || source.domain === 'lark') result.domain = source.domain;
if (typeof source.activated === 'boolean' || typeof source.activated === 'number') {
result.activated = source.activated;
}
return result;
}
function publicHealth(status, connected) {
if (connected) return { status: 'healthy', summary: '长连接运行正常', lastCheckedAt: Date.now() };
if (status?.configured === true) {
return { status: 'offline', summary: '机器人尚未连接', lastCheckedAt: Date.now() };
}
return { status: 'offline', summary: '尚未接入飞书机器人', lastCheckedAt: Date.now() };
}
function connectionState(status, registration, connected) {
if (connected) return 'connected';
if (status?.phase === 'error' || registration.state === 'error') return 'error';
if (status?.phase === 'connecting' || registration.state === 'saving') return 'connecting';
if (status?.phase === 'registering'
|| ['starting', 'qr_ready', 'polling', 'slow_down', 'domain_switched'].includes(registration.state)) {
return 'provisioning';
}
return 'disconnected';
}
async function qrCodeDataUrl(verificationUrl) {
return QRCode.toDataURL(verificationUrl, {
errorCorrectionLevel: 'M', margin: 1, width: 320, type: 'image/png',
});
}
async function publicProvisioning(registration, encodeQr) {
if (!registration.qrCodeUrl) return undefined;
return {
attemptId: String(registration.attempt),
verificationUrl: registration.qrCodeUrl,
qrCodeDataUrl: await encodeQr(registration.qrCodeUrl),
expiresAt: registration.expiresAt ?? Date.now() + (5 * 60_000),
pollIntervalMs: Math.max(800, Math.min(10_000, (registration.pollIntervalSeconds ?? 1.8) * 1000)),
};
}
function publicBotEntry(entry) {
const source = entry && typeof entry === 'object' ? entry : {};
if (!safeOpaqueId(source.botId)) return null;
const facts = connectionFacts(source.connection);
const connected = source.connected === true || facts.connected;
const registration = { state: 'idle' };
const result = {
botId: source.botId,
state: connectionState(source, registration, connected),
connected,
configured: source.configured === true,
bot: publicBot(source.bot),
health: publicHealth(source, connected),
};
const error = publicError(source.error);
if (error) result.error = error;
return result;
}
/** Exact redacted browser contract consumed by the Feishu settings client. */
export async function toPublicFeishuStatus(status, { encodeQr = qrCodeDataUrl } = {}) {
const source = status && typeof status === 'object' ? status : {};
const registration = publicRegistration(source.registration);
const facts = connectionFacts(source.connection);
const connected = source.connected === true || facts.connected;
const provisioning = await publicProvisioning(registration, encodeQr);
const error = publicError(source.error) ?? registration.error ?? null;
const bots = Array.isArray(source.bots)
? source.bots.map(publicBotEntry).filter(Boolean)
: [];
const snapshot = {
schemaVersion: source.schemaVersion === 2 ? 2 : 1,
revision: Number.isSafeInteger(source.revision) && source.revision >= 0 ? source.revision : 0,
state: connectionState(source, registration, connected),
connected,
configured: source.configured === true,
bot: publicBot(source.bot),
health: publicHealth(source, connected),
bots,
totals: {
configured: bots.length || (source.configured === true ? 1 : 0),
connected: bots.length ? bots.filter((bot) => bot.connected).length : (connected ? 1 : 0),
},
};
if (provisioning) snapshot.provisioning = provisioning;
if (error) snapshot.error = error;
return snapshot;
}
function badRequest(message) {
return { ok: false, error: { code: 'bad-request', message, details: { issues: [] } } };
}
function cancelled() {
return { ok: false, error: { code: 'cancelled', message: 'The Feishu request was cancelled.', details: {} } };
}
function internalFailure() {
return { ok: false, error: { code: 'internal', message: 'The Feishu integration operation failed.', details: {} } };
}
function validPayload(endpoint, payload) {
if (endpoint === FEISHU_ENDPOINTS.status) {
return hasOnlyKeys(payload, new Set()) ? null : 'This endpoint accepts an empty payload only.';
}
if (endpoint === FEISHU_ENDPOINTS.testConnection) {
return hasOnlyKeys(payload, new Set()) ? null : 'This endpoint accepts an empty payload only.';
}
if (endpoint === FEISHU_ENDPOINTS.beginProvisioning) {
if (!hasOnlyKeys(payload, new Set(['locale', 'replaceAttemptId']))) {
return 'Provisioning accepts locale and replaceAttemptId only.';
}
if (payload.locale !== undefined && payload.locale !== 'zh-CN') return 'The provisioning locale must be zh-CN.';
if (payload.replaceAttemptId !== undefined && !safeOpaqueId(payload.replaceAttemptId)) {
return 'replaceAttemptId must be a valid opaque id.';
}
return null;
}
if (endpoint === FEISHU_ENDPOINTS.pollProvisioning
|| endpoint === FEISHU_ENDPOINTS.cancelProvisioning) {
return hasOnlyKeys(payload, new Set(['attemptId'])) && safeOpaqueId(payload.attemptId)
? null
: 'A single valid attemptId is required.';
}
if (endpoint === FEISHU_ENDPOINTS.disconnect) {
return hasOnlyKeys(payload, new Set(['removeCredentials'])) && payload.removeCredentials === true
? null
: 'Disconnect requires removeCredentials=true.';
}
if (endpoint === FEISHU_MULTI_ENDPOINTS.reconnectBot
|| endpoint === FEISHU_MULTI_ENDPOINTS.disconnectBot) {
return hasOnlyKeys(payload, new Set(['botId'])) && safeOpaqueId(payload.botId)
? null
: 'A single valid botId is required.';
}
if (endpoint === FEISHU_MULTI_ENDPOINTS.deleteBot) {
return hasOnlyKeys(payload, new Set(['botId', 'confirm']))
&& safeOpaqueId(payload.botId) && payload.confirm === true
? null
: 'Deleting a bot requires a valid botId and confirm=true.';
}
return 'Unknown Feishu endpoint.';
}
function abortableDelay(milliseconds, signal) {
return new Promise((resolve, reject) => {
if (signal?.aborted) {
reject(signal.reason ?? new Error('aborted'));
return;
}
const timer = setTimeout(done, milliseconds);
timer.unref?.();
function done() {
signal?.removeEventListener('abort', aborted);
resolve();
}
function aborted() {
clearTimeout(timer);
reject(signal.reason ?? new Error('aborted'));
}
signal?.addEventListener('abort', aborted, { once: true });
});
}
async function statusForRegistration(controller, attemptId) {
if (typeof controller.registrationStatus === 'function') {
return controller.registrationStatus(attemptId);
}
return controller.status();
}
async function waitForQr(controller, initial, attemptId, signal) {
let current = initial;
const deadline = Date.now() + 15_000;
for (;;) {
const registration = publicRegistration(current?.registration);
if (registration.qrCodeUrl) return current;
if (['error', 'expired', 'cancelled'].includes(registration.state)) {
throw new Error('Provisioning stopped before the QR code was ready.');
}
if (Date.now() >= deadline) throw new Error('Provisioning QR code timed out.');
await abortableDelay(50, signal);
current = await statusForRegistration(controller, attemptId);
if (!current) throw new Error('The provisioning attempt is no longer active.');
}
}
function sameAttempt(status, attemptId) {
return String(publicRegistration(status?.registration).attempt) === attemptId;
}
function pollStatus(status) {
const registration = publicRegistration(status?.registration);
if (registration.state === 'succeeded') {
const connected = registration.botId
&& (status?.connected === true || connectionFacts(status?.connection).connected);
return connected ? 'connected' : 'connecting';
}
return POLL_STATUS_BY_REGISTRATION[registration.state] ?? 'failed';
}
function assertController(controller) {
if (!controller
|| typeof controller.status !== 'function'
|| typeof controller.startRegistration !== 'function'
|| typeof controller.cancelRegistration !== 'function'
|| typeof controller.disconnect !== 'function') {
throw new TypeError('A Feishu controller with status/start/cancel/disconnect is required');
}
}
/** DSH rc.6 handler: (endpoint, payload, signal) => Promise<RpcResult>. */
export function createFeishuRpcHandler(controller, { encodeQr = qrCodeDataUrl } = {}) {
assertController(controller);
const qrCache = new Map();
const attemptQr = new Map();
const cachedEncodeQr = (url) => {
let encoded = qrCache.get(url);
if (!encoded) {
if (qrCache.size >= 32) 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 cancelled();
if (!FEISHU_RPC_ENDPOINTS.includes(endpoint)) return badRequest('Unknown Feishu endpoint.');
const payloadFailure = validPayload(endpoint, payload);
if (payloadFailure) return badRequest(payloadFailure);
try {
let value;
if (endpoint === FEISHU_ENDPOINTS.status) {
value = await toPublicFeishuStatus(await controller.status(), { encodeQr: cachedEncodeQr });
} else if (endpoint === FEISHU_ENDPOINTS.beginProvisioning) {
if (payload.replaceAttemptId) {
await controller.cancelRegistration(payload.replaceAttemptId);
}
const started = await controller.startRegistration({ locale: payload.locale });
const attemptId = String(publicRegistration(started?.registration).attempt);
const ready = await waitForQr(controller, started, attemptId, signal);
value = (await toPublicFeishuStatus(ready, { encodeQr: cachedEncodeQr })).provisioning;
if (!value) throw new Error('Provisioning did not produce a QR code.');
attemptQr.set(attemptId, value.verificationUrl);
} else if (endpoint === FEISHU_ENDPOINTS.pollProvisioning) {
const current = await statusForRegistration(controller, payload.attemptId);
if (!current || !sameAttempt(current, payload.attemptId)) {
return badRequest('The provisioning attempt is no longer active.');
}
const registration = publicRegistration(current.registration);
const connection = await toPublicFeishuStatus(current, { encodeQr: cachedEncodeQr });
value = {
status: pollStatus(current),
...(registration.botId ? { botId: registration.botId } : {}),
...(connection.provisioning ? { provisioning: connection.provisioning } : {}),
...(registration.botId && connection.connected ? { connection } : {}),
...(connection.error ? { message: connection.error.message } : {}),
};
if (['connected', 'expired', 'failed'].includes(value.status)) {
const url = attemptQr.get(payload.attemptId);
if (url) qrCache.delete(url);
attemptQr.delete(payload.attemptId);
}
} else if (endpoint === FEISHU_ENDPOINTS.cancelProvisioning) {
const current = await statusForRegistration(controller, payload.attemptId);
if (!current || !sameAttempt(current, payload.attemptId)) {
return badRequest('The provisioning attempt is no longer active.');
}
const multi = typeof controller.registrationStatus === 'function';
const registration = publicRegistration(current.registration);
if (!multi && registration.state === 'saving') await controller.disconnect();
else await controller.cancelRegistration(payload.attemptId);
const url = attemptQr.get(payload.attemptId);
if (url) qrCache.delete(url);
attemptQr.delete(payload.attemptId);
value = { status: 'failed', message: 'Registration was cancelled.' };
} else if (endpoint === FEISHU_ENDPOINTS.testConnection) {
const current = await controller.status();
const alreadyConnected = current?.connected === true
|| connectionFacts(current?.connection).connected;
const checked = alreadyConnected || typeof controller.reconnect !== 'function'
? current
: await controller.reconnect();
value = await toPublicFeishuStatus(checked, { encodeQr: cachedEncodeQr });
} else if (endpoint === FEISHU_ENDPOINTS.disconnect) {
value = await toPublicFeishuStatus(await controller.disconnect(), { encodeQr: cachedEncodeQr });
} else if (endpoint === FEISHU_MULTI_ENDPOINTS.reconnectBot) {
if (typeof controller.reconnectBot !== 'function') throw new Error('Multi-bot reconnect is unavailable');
value = await toPublicFeishuStatus(await controller.reconnectBot(payload.botId), { encodeQr: cachedEncodeQr });
} else if (endpoint === FEISHU_MULTI_ENDPOINTS.disconnectBot) {
if (typeof controller.disconnectBot !== 'function') throw new Error('Multi-bot disconnect is unavailable');
value = await toPublicFeishuStatus(await controller.disconnectBot(payload.botId), { encodeQr: cachedEncodeQr });
} else {
if (typeof controller.deleteBot !== 'function') throw new Error('Multi-bot delete is unavailable');
value = await toPublicFeishuStatus(await controller.deleteBot(payload.botId), { encodeQr: cachedEncodeQr });
}
if (signal?.aborted) return cancelled();
return { ok: true, value };
} catch {
return signal?.aborted ? cancelled() : internalFailure();
}
};
}
/** Register the loopback-only `/feishu` logical channel. */
export function installFeishuRpc(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(
FEISHU_RPC_CHANNEL,
createFeishuRpcHandler(controller, options),
{ authority: 'loopback' },
);
}

View file

@ -0,0 +1,127 @@
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 delayMs = this.#retryDelays[Math.min(this.#retryIndex, this.#retryDelays.length - 1)];
this.#retryIndex += 1;
this.#logger.warn?.(
`[dsh-weixin] ${connected}/${configured} accounts connected; retrying in ${delayMs}ms`,
);
this.#schedule(delayMs);
return;
}
this.#retryIndex = 0;
this.#schedule(this.#healthyIntervalMs);
} catch (error) {
if (this.#closed) return;
const delayMs = this.#retryDelays[Math.min(this.#retryIndex, this.#retryDelays.length - 1)];
this.#retryIndex += 1;
this.#logger.warn?.(`[dsh-weixin] connection reconciliation failed; retrying in ${delayMs}ms`, error);
this.#schedule(delayMs);
}
}
}
export function createConnectionSupervisor(options) {
return new ConnectionSupervisor(options);
}

View file

@ -0,0 +1,32 @@
import { createProductionController } from './production.mjs';
import { installWeixinRpc } from './rpc.mjs';
export const name = 'dsh-weixin-host';
export const inject = ['connection', 'credentials', 'webServer'];
export async function apply(ctx, config = {}) {
if (config?.controller) return installWeixinRpc(ctx, config.controller, config.rpcOptions);
const production = await createProductionController(ctx, config, config.internals);
const disposeRpc = installWeixinRpc(ctx, production.controller, config.rpcOptions);
ctx.effect(() => async () => {
await production.close();
}, 'dsh-weixin: close account connections');
return disposeRpc;
}
export function createWeixinHostPlugin(config) {
return Object.freeze({ name, inject, apply: (ctx) => apply(ctx, config) });
}
export { createConnectionSupervisor, ConnectionSupervisor } from './connection-supervisor.mjs';
export { createProductionController } from './production.mjs';
export {
WEIXIN_ENDPOINTS,
WEIXIN_RPC_CHANNEL,
WEIXIN_RPC_ENDPOINTS,
createWeixinRpcHandler,
installWeixinRpc,
} from './rpc.mjs';
export { WeixinController } from '../../../../src/channels/weixin/weixin-controller.mjs';
export { WeixinRuntime } from '../../../../src/channels/weixin/weixin-runtime.mjs';

View file

@ -0,0 +1,119 @@
import { unlink } from 'node:fs/promises';
import { homedir } from 'node:os';
import { join, resolve } from 'node:path';
import { WeixinConfigStore } from '../../../../src/channels/weixin/config-store.mjs';
import { HarnessClient } from '../../../../src/channels/weixin/harness-client.mjs';
import { WeixinStateStore } from '../../../../src/channels/weixin/state-store.mjs';
import { createWeixinApi } from '../../../../src/channels/weixin/weixin-api.mjs';
import { WeixinController } from '../../../../src/channels/weixin/weixin-controller.mjs';
import { WeixinRuntime } from '../../../../src/channels/weixin/weixin-runtime.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-weixin 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-weixin'));
return {
root,
config: resolve(config.configPath ?? join(root, 'config.json')),
accounts: resolve(config.accountsDir ?? join(root, 'accounts')),
};
}
export async function createProductionController(ctx, config = {}, internals = {}) {
if (!ctx?.credentials) throw new TypeError('dsh-weixin requires ctx.credentials');
if (!ctx?.webServer) throw new TypeError('dsh-weixin requires ctx.webServer');
const ConfigStore = internals.ConfigStore ?? WeixinConfigStore;
const StateStore = internals.StateStore ?? WeixinStateStore;
const Harness = internals.HarnessClient ?? HarnessClient;
const Controller = internals.Controller ?? WeixinController;
const Runtime = internals.Runtime ?? WeixinRuntime;
const api = internals.api ?? createWeixinApi();
const createSupervisor = internals.createConnectionSupervisor ?? createConnectionSupervisor;
const logger = typeof ctx.logger === 'function'
? ctx.logger('dsh-weixin')
: (ctx.logger ?? console);
const paths = pluginPaths(config);
const configStore = await new ConfigStore(paths.config).load();
const stateStores = new Map();
const statePath = (botId) => resolve(paths.accounts, 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({
api,
credentials: ctx.credentials,
configStore,
logger,
createRuntime: async ({ botId, config: accountConfig, token }) => {
const state = await stateFor(botId);
return new Runtime({
api,
config: accountConfig,
token,
harness,
state,
replyTimeoutMs: config.replyTimeoutMs ?? 600_000,
maxMessageChars: config.maxMessageChars ?? 4_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') {
await state.remove();
return;
}
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,177 @@
import QRCode from 'qrcode';
export const WEIXIN_RPC_CHANNEL = '/weixin';
export const WEIXIN_ENDPOINTS = Object.freeze({
status: 'connection.status',
beginProvisioning: 'provision.begin',
pollProvisioning: 'provision.poll',
submitVerification: 'provision.verify',
cancelProvisioning: 'provision.cancel',
reconnectBot: 'bot.reconnect',
deleteBot: 'bot.delete',
});
export const WEIXIN_RPC_ENDPOINTS = Object.freeze(Object.values(WEIXIN_ENDPOINTS));
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 === WEIXIN_ENDPOINTS.status) {
return exactKeys(payload, []) ? null : 'connection.status does not accept fields.';
}
if (endpoint === WEIXIN_ENDPOINTS.beginProvisioning) {
return exactKeys(payload, ['locale']) && (payload.locale === undefined || payload.locale === 'zh-CN')
? null
: 'provision.begin received unsupported fields.';
}
if ([WEIXIN_ENDPOINTS.pollProvisioning, WEIXIN_ENDPOINTS.cancelProvisioning].includes(endpoint)) {
return exactKeys(payload, ['attemptId']) && validId(payload.attemptId)
? null
: `${endpoint} requires an attemptId.`;
}
if (endpoint === WEIXIN_ENDPOINTS.submitVerification) {
return exactKeys(payload, ['attemptId', 'verifyCode'])
&& validId(payload.attemptId)
&& typeof payload.verifyCode === 'string'
&& /^\d{4,8}$/.test(payload.verifyCode)
? null
: 'provision.verify requires an attemptId and a 4-to-8-digit code.';
}
if (endpoint === WEIXIN_ENDPOINTS.reconnectBot) {
return exactKeys(payload, ['botId']) && validId(payload.botId)
? null
: 'bot.reconnect requires a botId.';
}
if (endpoint === WEIXIN_ENDPOINTS.deleteBot) {
return exactKeys(payload, ['botId', 'confirm']) && validId(payload.botId) && payload.confirm === true
? null
: 'bot.delete requires a botId and confirm=true.';
}
return 'Unknown Weixin endpoint.';
}
function badRequest(message) {
return { ok: false, error: { code: 'bad-request', message } };
}
function cancelled() {
return { ok: false, error: { code: 'cancelled', message: 'The request was cancelled.' } };
}
function internalFailure() {
return {
ok: false,
error: { code: 'weixin-operation-failed', message: '微信操作失败,请稍后重试。' },
};
}
async function qrDataUrl(value) {
return QRCode.toDataURL(value, {
type: 'image/png',
errorCorrectionLevel: 'M',
margin: 2,
width: 320,
});
}
async function withEncodedQr(value, encodeQr) {
if (!value || !value.verificationUrl) return value;
return {
...value,
qrCodeDataUrl: await encodeQr(value.verificationUrl),
};
}
async function publicStatus(status, encodeQr) {
const safe = structuredClone(status);
if (safe.provisioning) safe.provisioning = await withEncodedQr(safe.provisioning, encodeQr);
return safe;
}
function assertController(controller) {
if (!controller
|| typeof controller.status !== 'function'
|| typeof controller.startProvisioning !== 'function'
|| typeof controller.registrationStatus !== 'function'
|| typeof controller.submitVerification !== 'function'
|| typeof controller.cancelProvisioning !== 'function'
|| typeof controller.reconnectBot !== 'function'
|| typeof controller.deleteBot !== 'function') {
throw new TypeError('A complete Weixin controller is required');
}
}
export function createWeixinRpcHandler(controller, { encodeQr = qrDataUrl } = {}) {
assertController(controller);
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 cancelled();
if (!WEIXIN_RPC_ENDPOINTS.includes(endpoint)) return badRequest('Unknown Weixin endpoint.');
const invalid = payloadFailure(endpoint, payload);
if (invalid) return badRequest(invalid);
try {
let value;
if (endpoint === WEIXIN_ENDPOINTS.status) {
value = await publicStatus(await controller.status(), cachedEncode);
} else if (endpoint === WEIXIN_ENDPOINTS.beginProvisioning) {
const started = await controller.startProvisioning();
if (signal?.aborted) {
await controller.cancelProvisioning(started.attemptId);
return cancelled();
}
value = await withEncodedQr(started, cachedEncode);
} else if (endpoint === WEIXIN_ENDPOINTS.pollProvisioning) {
const current = await controller.registrationStatus(payload.attemptId);
if (!current) return badRequest('The provisioning attempt no longer exists.');
value = await withEncodedQr(current, cachedEncode);
} else if (endpoint === WEIXIN_ENDPOINTS.submitVerification) {
value = await withEncodedQr(
await controller.submitVerification(payload.attemptId, payload.verifyCode),
cachedEncode,
);
} else if (endpoint === WEIXIN_ENDPOINTS.cancelProvisioning) {
value = await controller.cancelProvisioning(payload.attemptId);
if (!value) return badRequest('The provisioning attempt no longer exists.');
} else if (endpoint === WEIXIN_ENDPOINTS.reconnectBot) {
value = await publicStatus(await controller.reconnectBot(payload.botId), cachedEncode);
} else {
value = await publicStatus(await controller.deleteBot(payload.botId), cachedEncode);
}
return signal?.aborted ? cancelled() : { ok: true, value };
} catch {
return signal?.aborted ? cancelled() : internalFailure();
}
};
}
export function installWeixinRpc(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(
WEIXIN_RPC_CHANNEL,
createWeixinRpcHandler(controller, options),
{ authority: 'loopback' },
);
}

View file

@ -1,6 +1,6 @@
import { apply as applyDingtalk } from '@xmanrui/dsh-dingtalk';
import { apply as applyFeishu } from '@xmanrui/dsh-feishu';
import { apply as applyWeixin } from '@xmanrui/dsh-weixin';
import { apply as applyDingtalk } from './channels/dingtalk/index.mjs';
import { apply as applyFeishu } from './channels/feishu/index.mjs';
import { apply as applyWeixin } from './channels/weixin/index.mjs';
export const name = 'dsh-im-host';
export const inject = ['connection', 'credentials', 'webServer'];