dsh-im-ops/plugin-src/host/index.mjs
oliver 41f13bb90c Expose resolveSessionPeer on ctx.dshIm for ops schedulers (ops.23).
WhatsApp registers a peer resolver so cron_create can default IM delivery to the calling session's bot/target.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-05 21:48:00 +08:00

162 lines
6.5 KiB
JavaScript

import { apply as applyDingtalk } from './channels/dingtalk/index.mjs';
import { apply as applyDiscord } from './channels/discord/index.mjs';
import { apply as applyOffice } from './channels/office/index.mjs';
import { apply as applyFeishu } from './channels/feishu/index.mjs';
import { apply as applyQq } from './channels/qq/index.mjs';
import { apply as applySlack } from './channels/slack/index.mjs';
import { apply as applyTelegram } from './channels/telegram/index.mjs';
import { apply as applyWecom } from './channels/wecom/index.mjs';
import { apply as applyWeixin } from './channels/weixin/index.mjs';
import { apply as applyWhatsapp } from './channels/whatsapp/index.mjs';
import { installOutboundArtifactTool } from '../../src/channels/shared/semantic/artifact.mjs';
import { setImHostLanguage } from '../../src/channels/shared/i18n.mjs';
import { installDeliveryRpc } from './delivery-rpc.mjs';
import { installDeliveryHttp } from './delivery-http.mjs';
import { createDeliveryService } from './delivery-service.mjs';
import { installUpdateRpc } from './update-rpc.mjs';
export const name = 'dsh-im-host';
export const inject = [
'connection',
'credentials',
'typertGateway',
];
function channelConfig(config, name, deliveryService) {
const channel = config[name] ?? {};
const withAuthority = config.rpcAuthority === undefined
? channel
: { ...channel, rpcAuthority: config.rpcAuthority };
return name === 'office' ? withAuthority : { ...withAuthority, deliveryService };
}
export function createImHostPlugin(internals = {}) {
const startUpdate = internals.installUpdateRpc ?? installUpdateRpc;
const startDelivery = internals.installDeliveryRpc ?? installDeliveryRpc;
const startDeliveryHttp = internals.installDeliveryHttp ?? installDeliveryHttp;
const makeDeliveryService = internals.createDeliveryService ?? createDeliveryService;
const startFeishu = internals.applyFeishu ?? applyFeishu;
const startWeixin = internals.applyWeixin ?? applyWeixin;
const startDingtalk = internals.applyDingtalk ?? applyDingtalk;
const startWecom = internals.applyWecom ?? applyWecom;
const startQq = internals.applyQq ?? applyQq;
const startSlack = internals.applySlack ?? applySlack;
const startTelegram = internals.applyTelegram ?? applyTelegram;
const startDiscord = internals.applyDiscord ?? applyDiscord;
const startOffice = internals.applyOffice ?? applyOffice;
const startWhatsapp = internals.applyWhatsapp ?? applyWhatsapp;
const channels = [
['feishu', startFeishu],
['weixin', startWeixin],
['dingtalk', startDingtalk],
['wecom', startWecom],
['qq', startQq],
['slack', startSlack],
['telegram', startTelegram],
['discord', startDiscord],
['whatsapp', startWhatsapp],
['office', startOffice],
];
return Object.freeze({
name,
inject,
async apply(ctx, config = {}) {
const deliveryService = makeDeliveryService();
const peerResolvers = [];
if (typeof ctx?.provide === 'function') {
ctx.provide('dshIm', Object.freeze({
send: (botId, targetId, text, options) => (
deliveryService.send(botId, targetId, text, options)
),
listTargets: async (botId) => (await deliveryService.listTargets(botId)).targets,
/**
* Resolve the IM channel peer for a Harness session (botId + conversationKey).
* Used by ops schedulers to default proactive delivery back to the caller.
*/
resolveSessionPeer: async (sessionId) => {
const id = typeof sessionId === 'string' ? sessionId.trim() : '';
if (!id) return null;
for (const resolve of peerResolvers) {
try {
const peer = await resolve(id);
if (peer && typeof peer.botId === 'string' && peer.botId.trim()) return peer;
} catch {
// try next channel
}
}
return null;
},
/** @param {(sessionId: string) => Promise<object|null>} resolve */
registerPeerResolver: (resolve) => {
if (typeof resolve !== 'function') return () => {};
peerResolvers.push(resolve);
return () => {
const index = peerResolvers.indexOf(resolve);
if (index >= 0) peerResolvers.splice(index, 1);
};
},
}));
}
const activate = async (readyCtx) => {
await activateChannels(readyCtx, config, deliveryService);
};
if (typeof ctx?.inject === 'function') {
const modern = typeof ctx?.typertGateway?.stream === 'function';
await ctx.inject(
modern ? ['sessionController', 'workspaceController'] : ['apiProxy'],
activate,
);
ctx.inject(['webServer'], (httpCtx) => {
startDeliveryHttp(httpCtx, deliveryService);
});
return;
}
await activate(ctx);
if (ctx?.webServer?.register && typeof ctx?.effect === 'function') {
startDeliveryHttp(ctx, deliveryService);
}
},
});
async function activateChannels(ctx, config, deliveryService) {
setImHostLanguage(config.language ?? process.env.DSH_IM_LANGUAGE);
if (typeof ctx?.inject === 'function') {
ctx.inject(['tools', 'systemPrompt'], (artifactCtx) => {
installOutboundArtifactTool(artifactCtx);
});
} else {
installOutboundArtifactTool(ctx);
}
const logger = typeof ctx?.logger === 'function'
? ctx.logger(name)
: (ctx?.logger ?? console);
if (ctx?.connection?.rpc) {
try {
startUpdate(ctx);
} catch (error) {
logger.error?.('[dsh-im] failed to activate update management; continuing with channels', error);
}
try {
startDelivery(ctx, deliveryService, { authority: config.rpcAuthority });
} catch (error) {
logger.error?.('[dsh-im] failed to activate delivery management; continuing with channels', error);
}
}
const failures = [];
for (const [channel, start] of channels) {
try {
await start(ctx, channelConfig(config, channel, deliveryService));
} catch (error) {
failures.push(error);
logger.error?.(`[dsh-im] failed to activate ${channel}; continuing with the remaining channels`, error);
}
}
if (failures.length === channels.length) {
throw new AggregateError(failures, 'dsh-im failed to activate every channel');
}
}
}
export async function apply(ctx, config = {}) {
return createImHostPlugin().apply(ctx, config);
}