feat: add stable proactive delivery

This commit is contained in:
xmanrui 2026-08-30 11:24:53 +08:00
parent 410cb822d7
commit 31792d7bc8
88 changed files with 6427 additions and 488 deletions

View file

@ -10,6 +10,8 @@ export async function apply(ctx, config = {}) {
}
const production = await createProductionController(ctx, config, config.internals);
const unregisterDelivery = config.deliveryService && production.deliveryAdapter
? config.deliveryService.registerAdapter(production.deliveryAdapter) : undefined;
const disposeRpc = installDingtalkRpc(
ctx,
production.controller,
@ -17,6 +19,7 @@ export async function apply(ctx, config = {}) {
config.rpcAuthority,
);
ctx.effect(() => async () => {
await unregisterDelivery?.();
await production.close();
}, 'dsh-dingtalk: close bot connections');
return disposeRpc;

View file

@ -15,6 +15,7 @@ import {
observeBotWorkspaceRemovals,
} from '../../../../src/channels/shared/bot-workspace-store.mjs';
import { listAgentPresetCatalog } from '../../../../src/channels/shared/agent-preset.mjs';
import { createDeliveryAdapter } from '../../delivery-adapter.mjs';
import { createConnectionSupervisor } from './connection-supervisor.mjs';
import { createHarnessCommandExecutor } from '../../harness-command-executor.mjs';
import { harnessConnection } from '../../harness-connection.mjs';
@ -149,6 +150,9 @@ export async function createProductionController(ctx, config = {}, internals = {
}).start();
return {
controller,
deliveryAdapter: createDeliveryAdapter({
channel: 'dingtalk', workspaces, coreController, stateFor,
}),
ready: supervisor.ready,
async close() {
await supervisor.close();

View file

@ -9,8 +9,13 @@ export async function apply(ctx, config = {}) {
return installDiscordRpc(ctx, config.controller, config.rpcAuthority);
}
const production = await createProductionController(ctx, config, config.internals ?? {});
const unregisterDelivery = config.deliveryService && production.deliveryAdapter
? config.deliveryService.registerAdapter(production.deliveryAdapter) : undefined;
const disposeRpc = installDiscordRpc(ctx, production.controller, config.rpcAuthority);
ctx.effect(() => async () => production.close(), 'dsh-im: close Discord bot connections');
ctx.effect(() => async () => {
await unregisterDelivery?.();
await production.close();
}, 'dsh-im: close Discord bot connections');
return disposeRpc;
}

View file

@ -29,6 +29,8 @@ export async function apply(ctx, config = {}) {
}
const production = await createProductionController(ctx, config);
const unregisterDelivery = config.deliveryService && production.deliveryAdapter
? config.deliveryService.registerAdapter(production.deliveryAdapter) : undefined;
const disposeRpc = installFeishuRpc(
ctx,
production.controller,
@ -36,6 +38,7 @@ export async function apply(ctx, config = {}) {
config.rpcAuthority,
);
ctx.effect(() => async () => {
await unregisterDelivery?.();
await production.close();
}, 'dsh-feishu: close controller and live connection');
return disposeRpc;

View file

@ -23,6 +23,7 @@ import {
observeBotWorkspaceRemovals,
} from '../../../../src/channels/shared/bot-workspace-store.mjs';
import { listAgentPresetCatalog } from '../../../../src/channels/shared/agent-preset.mjs';
import { createDeliveryAdapter } from '../../delivery-adapter.mjs';
function webSocketProxyUrl(env) {
for (const key of ['https_proxy', 'HTTPS_PROXY', 'http_proxy', 'HTTP_PROXY']) {
@ -204,6 +205,9 @@ export async function createProductionController(ctx, config = {}, internals = {
}).start();
return {
controller,
deliveryAdapter: createDeliveryAdapter({
channel: 'feishu', workspaces, coreController, stateFor: stateForBotId,
}),
ready: supervisor.ready,
async close() {
await supervisor.close();

View file

@ -9,13 +9,18 @@ export async function apply(ctx, config = {}) {
return installQqRpc(ctx, config.controller, config.rpcOptions, config.rpcAuthority);
}
const production = await createProductionController(ctx, config, config.internals);
const unregisterDelivery = config.deliveryService && production.deliveryAdapter
? config.deliveryService.registerAdapter(production.deliveryAdapter) : undefined;
const disposeRpc = installQqRpc(
ctx,
production.controller,
config.rpcOptions,
config.rpcAuthority,
);
ctx.effect(() => async () => production.close(), 'dsh-im: close QQ bot connections');
ctx.effect(() => async () => {
await unregisterDelivery?.();
await production.close();
}, 'dsh-im: close QQ bot connections');
return disposeRpc;
}

View file

@ -15,6 +15,7 @@ import {
observeBotWorkspaceRemovals,
} from '../../../../src/channels/shared/bot-workspace-store.mjs';
import { listAgentPresetCatalog } from '../../../../src/channels/shared/agent-preset.mjs';
import { createDeliveryAdapter } from '../../delivery-adapter.mjs';
import { createConnectionSupervisor } from './connection-supervisor.mjs';
import { createHarnessCommandExecutor } from '../../harness-command-executor.mjs';
import { harnessConnection } from '../../harness-connection.mjs';
@ -139,6 +140,7 @@ export async function createProductionController(ctx, config = {}, internals = {
}).start();
return {
controller,
deliveryAdapter: createDeliveryAdapter({ channel: 'qq', workspaces, coreController, stateFor }),
ready: supervisor.ready,
async close() {
await supervisor.close();

View file

@ -13,6 +13,10 @@ import {
observeBotWorkspaceRemovals,
} from '../../../../src/channels/shared/bot-workspace-store.mjs';
import { listAgentPresetCatalog } from '../../../../src/channels/shared/agent-preset.mjs';
import {
createDeliveryAdapter,
supportsDeliveryChannel,
} from '../../delivery-adapter.mjs';
export function pluginPaths(config, channel) {
const dshHome = resolve(config.dshHome ?? process.env.DSH_HOME ?? join(homedir(), '.dsh'));
@ -142,6 +146,9 @@ export async function createTokenProductionController(ctx, config, internals, de
}).start();
return {
controller,
...(supportsDeliveryChannel(channel) ? {
deliveryAdapter: createDeliveryAdapter({ channel, workspaces, coreController, stateFor }),
} : {}),
ready: supervisor.ready,
async close() {
await supervisor.close();

View file

@ -7,8 +7,13 @@ export const inject = ['connection', 'credentials', 'typertGateway'];
export async function apply(ctx, config = {}) {
if (config?.controller) return installSlackRpc(ctx, config.controller, config.rpcAuthority);
const production = await createProductionController(ctx, config, config.internals ?? {});
const unregisterDelivery = config.deliveryService && production.deliveryAdapter
? config.deliveryService.registerAdapter(production.deliveryAdapter) : undefined;
const disposeRpc = installSlackRpc(ctx, production.controller, config.rpcAuthority);
ctx.effect(() => async () => production.close(), 'dsh-im: close Slack bot connections');
ctx.effect(() => async () => {
await unregisterDelivery?.();
await production.close();
}, 'dsh-im: close Slack bot connections');
return disposeRpc;
}

View file

@ -13,6 +13,7 @@ import {
observeBotWorkspaceRemovals,
} from '../../../../src/channels/shared/bot-workspace-store.mjs';
import { listAgentPresetCatalog } from '../../../../src/channels/shared/agent-preset.mjs';
import { createDeliveryAdapter } from '../../delivery-adapter.mjs';
import { createTokenConnectionSupervisor } from '../shared/connection-supervisor.mjs';
import { pluginPaths } from '../shared/production.mjs';
import { createHarnessCommandExecutor } from '../../harness-command-executor.mjs';
@ -129,6 +130,9 @@ export async function createProductionController(ctx, config = {}, internals = {
}).start();
return {
controller,
deliveryAdapter: createDeliveryAdapter({
channel: 'slack', workspaces, coreController, stateFor,
}),
ready: supervisor.ready,
async close() {
await supervisor.close();

View file

@ -9,8 +9,13 @@ export async function apply(ctx, config = {}) {
return installTelegramRpc(ctx, config.controller, config.rpcAuthority);
}
const production = await createProductionController(ctx, config, config.internals ?? {});
const unregisterDelivery = config.deliveryService && production.deliveryAdapter
? config.deliveryService.registerAdapter(production.deliveryAdapter) : undefined;
const disposeRpc = installTelegramRpc(ctx, production.controller, config.rpcAuthority);
ctx.effect(() => async () => production.close(), 'dsh-im: close Telegram bot connections');
ctx.effect(() => async () => {
await unregisterDelivery?.();
await production.close();
}, 'dsh-im: close Telegram bot connections');
return disposeRpc;
}

View file

@ -9,13 +9,18 @@ export async function apply(ctx, config = {}) {
return installWecomRpc(ctx, config.controller, config.rpcOptions, config.rpcAuthority);
}
const production = await createProductionController(ctx, config, config.internals);
const unregisterDelivery = config.deliveryService && production.deliveryAdapter
? config.deliveryService.registerAdapter(production.deliveryAdapter) : undefined;
const disposeRpc = installWecomRpc(
ctx,
production.controller,
config.rpcOptions,
config.rpcAuthority,
);
ctx.effect(() => async () => production.close(), 'dsh-im: close Enterprise WeChat bot connections');
ctx.effect(() => async () => {
await unregisterDelivery?.();
await production.close();
}, 'dsh-im: close Enterprise WeChat bot connections');
return disposeRpc;
}

View file

@ -15,6 +15,7 @@ import {
observeBotWorkspaceRemovals,
} from '../../../../src/channels/shared/bot-workspace-store.mjs';
import { listAgentPresetCatalog } from '../../../../src/channels/shared/agent-preset.mjs';
import { createDeliveryAdapter } from '../../delivery-adapter.mjs';
import { createConnectionSupervisor } from './connection-supervisor.mjs';
import { createHarnessCommandExecutor } from '../../harness-command-executor.mjs';
import { harnessConnection } from '../../harness-connection.mjs';
@ -143,6 +144,9 @@ export async function createProductionController(ctx, config = {}, internals = {
}).start();
return {
controller,
deliveryAdapter: createDeliveryAdapter({
channel: 'wecom', workspaces, coreController, stateFor,
}),
ready: supervisor.ready,
async close() {
await supervisor.close();

View file

@ -10,6 +10,8 @@ export async function apply(ctx, config = {}) {
}
const production = await createProductionController(ctx, config, config.internals);
const unregisterDelivery = config.deliveryService && production.deliveryAdapter
? config.deliveryService.registerAdapter(production.deliveryAdapter) : undefined;
const disposeRpc = installWeixinRpc(
ctx,
production.controller,
@ -17,6 +19,7 @@ export async function apply(ctx, config = {}) {
config.rpcAuthority,
);
ctx.effect(() => async () => {
await unregisterDelivery?.();
await production.close();
}, 'dsh-weixin: close account connections');
return disposeRpc;

View file

@ -18,6 +18,7 @@ import {
observeBotWorkspaceRemovals,
} from '../../../../src/channels/shared/bot-workspace-store.mjs';
import { listAgentPresetCatalog } from '../../../../src/channels/shared/agent-preset.mjs';
import { createDeliveryAdapter } from '../../delivery-adapter.mjs';
import { createConnectionSupervisor } from './connection-supervisor.mjs';
import { createHarnessCommandExecutor } from '../../harness-command-executor.mjs';
import { harnessConnection } from '../../harness-connection.mjs';
@ -146,6 +147,9 @@ export async function createProductionController(ctx, config = {}, internals = {
}).start();
return {
controller,
deliveryAdapter: createDeliveryAdapter({
channel: 'weixin', workspaces, coreController, stateFor,
}),
ready: supervisor.ready,
async close() {
await supervisor.close();

View file

@ -9,13 +9,18 @@ export async function apply(ctx, config = {}) {
return installWhatsappRpc(ctx, config.controller, config.rpcOptions, config.rpcAuthority);
}
const production = await createProductionController(ctx, config, config.internals ?? {});
const unregisterDelivery = config.deliveryService && production.deliveryAdapter
? config.deliveryService.registerAdapter(production.deliveryAdapter) : undefined;
const disposeRpc = installWhatsappRpc(
ctx,
production.controller,
config.rpcOptions,
config.rpcAuthority,
);
ctx.effect(() => async () => production.close(), 'dsh-im: close WhatsApp Web connections');
ctx.effect(() => async () => {
await unregisterDelivery?.();
await production.close();
}, 'dsh-im: close WhatsApp Web connections');
return disposeRpc;
}

View file

@ -15,6 +15,7 @@ import {
observeBotWorkspaceRemovals,
} from '../../../../src/channels/shared/bot-workspace-store.mjs';
import { listAgentPresetCatalog } from '../../../../src/channels/shared/agent-preset.mjs';
import { createDeliveryAdapter } from '../../delivery-adapter.mjs';
import { createTokenConnectionSupervisor } from '../shared/connection-supervisor.mjs';
import { createHarnessCommandExecutor } from '../../harness-command-executor.mjs';
import { harnessConnection } from '../../harness-connection.mjs';
@ -151,6 +152,9 @@ export async function createProductionController(ctx, config = {}, internals = {
}).start();
return {
controller,
deliveryAdapter: createDeliveryAdapter({
channel: 'whatsapp', workspaces, coreController, stateFor,
}),
ready: supervisor.ready,
async close() {
await supervisor.close();

View file

@ -0,0 +1,179 @@
import { deliverySuggestionsFromSessions } from './delivery-suggestions.mjs';
const CHANNELS = new Set([
'weixin',
'feishu',
'dingtalk',
'wecom',
'qq',
'slack',
'telegram',
'discord',
'whatsapp',
]);
export function supportsDeliveryChannel(channel) {
return CHANNELS.has(channel);
}
function invalidTarget(message) {
const error = new Error(message);
error.code = 'invalid-target';
return error;
}
function isRecord(value) {
return value !== null && typeof value === 'object' && !Array.isArray(value);
}
function exactKeys(value, allowed, label) {
if (!isRecord(value)) throw invalidTarget(`${label} must be an object`);
const extra = Object.keys(value).find((key) => !allowed.includes(key));
if (extra !== undefined) throw invalidTarget(`${label} contains an unknown field`);
}
function requiredString(value, label) {
if (typeof value !== 'string' || !value || value.trim() !== value) {
throw invalidTarget(`${label} must be a non-empty string without surrounding whitespace`);
}
return value;
}
function targetName(value) {
if (typeof value !== 'string' || !value.trim() || value.trim().length > 80) {
throw invalidTarget('target.name must contain 1 to 80 characters');
}
return value.trim();
}
function routeWithStrings(route, fields) {
exactKeys(route, fields, 'route');
const normalized = {};
for (const field of fields) normalized[field] = requiredString(route[field], `route.${field}`);
return normalized;
}
function oneOf(value, choices) {
if (!choices.includes(value)) throw invalidTarget('target kind is not supported by this channel');
return value;
}
function normalizeRoute(channel, kind, route) {
switch (channel) {
case 'weixin':
oneOf(kind, ['user']);
return routeWithStrings(route, ['toUserId']);
case 'feishu':
oneOf(kind, ['user', 'group']);
return routeWithStrings(route, kind === 'user' ? ['openId'] : ['chatId']);
case 'dingtalk':
oneOf(kind, ['user', 'group']);
return routeWithStrings(route, kind === 'user' ? ['userId'] : ['openConversationId']);
case 'wecom':
oneOf(kind, ['user', 'group']);
return routeWithStrings(route, ['chatId']);
case 'qq':
oneOf(kind, ['user', 'group']);
return routeWithStrings(route, kind === 'user' ? ['userOpenId'] : ['groupOpenId']);
case 'slack': {
oneOf(kind, ['conversation', 'thread']);
const normalized = routeWithStrings(
route,
kind === 'conversation' ? ['channelId'] : ['channelId', 'threadTs'],
);
return normalized;
}
case 'telegram': {
oneOf(kind, ['chat', 'topic']);
exactKeys(route, kind === 'chat' ? ['chatId'] : ['chatId', 'messageThreadId'], 'route');
const chatId = requiredString(route.chatId, 'route.chatId');
if (!/^-?\d+$/.test(chatId)) throw invalidTarget('route.chatId must be a decimal string');
if (kind === 'chat') return { chatId };
if (!Number.isSafeInteger(route.messageThreadId) || route.messageThreadId <= 0) {
throw invalidTarget('route.messageThreadId must be a positive integer');
}
return { chatId, messageThreadId: route.messageThreadId };
}
case 'discord':
oneOf(kind, ['channel']);
return routeWithStrings(route, ['channelId']);
case 'whatsapp': {
oneOf(kind, ['user', 'group']);
const normalized = routeWithStrings(route, ['jid']);
const valid = kind === 'user'
? /^\d{5,32}@(s\.whatsapp\.net|lid)$/.test(normalized.jid)
: /^\d{5,32}(?:-\d{1,32})?@g\.us$/.test(normalized.jid);
if (!valid) {
throw invalidTarget(`route.jid must be a ${kind} WhatsApp JID`);
}
return normalized;
}
default:
throw new TypeError(`Unsupported delivery channel: ${channel}`);
}
}
/** Strictly validate and copy one channel delivery target. */
export function normalizeDeliveryTarget(channel, value, { targetIdRequired = true } = {}) {
if (!supportsDeliveryChannel(channel)) throw new TypeError(`Unsupported delivery channel: ${channel}`);
exactKeys(
value,
targetIdRequired ? ['targetId', 'name', 'kind', 'route'] : ['name', 'kind', 'route'],
'target',
);
const normalized = {};
if (targetIdRequired) normalized.targetId = requiredString(value.targetId, 'target.targetId');
if (value.name !== undefined) normalized.name = targetName(value.name);
normalized.kind = requiredString(value.kind, 'target.kind');
normalized.route = normalizeRoute(channel, normalized.kind, value.route);
return normalized;
}
/** Bind one channel's existing workspace store and unwrapped controller to DeliveryService. */
export function createDeliveryAdapter({ channel, workspaces, coreController, stateFor }) {
if (!supportsDeliveryChannel(channel)) throw new TypeError(`Unsupported delivery channel: ${channel}`);
if (!workspaces || typeof workspaces !== 'object') {
throw new TypeError('delivery adapter requires a workspace store');
}
if (!coreController || typeof coreController !== 'object') {
throw new TypeError('delivery adapter requires a proactive text controller');
}
if (typeof stateFor !== 'function') {
throw new TypeError('delivery adapter requires a bot state getter');
}
return Object.freeze({
channel,
ownsBot: (botId) => workspaces.has(botId),
listTargets: (botId) => workspaces.listDeliveryTargets(botId),
async listSuggestions(botId) {
const state = await stateFor(botId);
if (!state || typeof state.snapshot !== 'function') {
throw new TypeError('delivery suggestion state cannot be inspected');
}
const suggestions = deliverySuggestionsFromSessions(channel, state.snapshot()?.sessions);
return suggestions.map((suggestion) => normalizeDeliveryTarget(
channel,
suggestion,
{ targetIdRequired: false },
));
},
createTarget: (botId, target) => workspaces.createDeliveryTarget(
botId,
normalizeDeliveryTarget(channel, target),
),
updateTarget: (botId, targetId, replacement) => workspaces.updateDeliveryTarget(
botId,
targetId,
normalizeDeliveryTarget(channel, replacement, { targetIdRequired: false }),
),
deleteTarget: (botId, targetId) => workspaces.deleteDeliveryTarget(botId, targetId),
async sendText(botId, target, text, options = {}) {
const normalized = normalizeDeliveryTarget(channel, target);
if (typeof coreController.sendProactiveText !== 'function') {
throw new TypeError('delivery controller cannot send proactive text');
}
await coreController.sendProactiveText(botId, normalized, text, options);
return { sent: true };
},
});
}

View file

@ -0,0 +1,158 @@
import { resolveRpcAuthority } from './rpc-authority.mjs';
export const DELIVERY_RPC_CHANNEL = '/dsh-im-delivery';
export const DELIVERY_TEST_MESSAGE = 'DSH-IM 主动投递测试成功。';
export const DELIVERY_ENDPOINTS = Object.freeze({
send: 'message.send',
listTargets: 'target.list',
listSuggestions: 'target.suggestion.list',
createTarget: 'target.create',
updateTarget: 'target.update',
deleteTarget: 'target.delete',
testTarget: 'target.test',
});
const ENDPOINTS = new Set(Object.values(DELIVERY_ENDPOINTS));
const PUBLIC_ERRORS = new Set([
'bad-request',
'unknown-bot',
'unknown-target',
'target-conflict',
'invalid-target',
'bot-not-connected',
'target-rejected',
'delivery-failed',
'cancelled',
]);
function isRecord(value) {
return value !== null && typeof value === 'object' && !Array.isArray(value);
}
function exactKeys(value, keys) {
return isRecord(value)
&& Object.keys(value).length === keys.length
&& Object.keys(value).every((key) => keys.includes(key));
}
function validBotId(value) {
return typeof value === 'string' && /^[A-Za-z0-9_-]{1,128}$/.test(value);
}
function validTargetId(value) {
return typeof value === 'string' && /^[A-Za-z0-9._:@-]{1,128}$/.test(value);
}
function validTarget(value, { targetId }) {
const keys = targetId ? ['targetId', 'name', 'kind', 'route'] : ['name', 'kind', 'route'];
if (!isRecord(value) || Object.keys(value).some((key) => !keys.includes(key))) return false;
if (targetId && !validTargetId(value.targetId)) return false;
if (value.name !== undefined
&& (typeof value.name !== 'string' || !value.name.trim() || value.name.trim().length > 80)) return false;
return typeof value.kind === 'string' && /^[a-z][a-z0-9-]{0,31}$/.test(value.kind)
&& isRecord(value.route);
}
function validDraftTarget(value) {
return exactKeys(value, ['kind', 'route'])
&& typeof value.kind === 'string' && /^[a-z][a-z0-9-]{0,31}$/.test(value.kind)
&& isRecord(value.route);
}
function validPayload(endpoint, payload) {
if (!ENDPOINTS.has(endpoint) || !isRecord(payload)) return false;
if (endpoint === DELIVERY_ENDPOINTS.send) {
return exactKeys(payload, ['botId', 'targetId', 'text'])
&& validBotId(payload.botId) && validTargetId(payload.targetId)
&& typeof payload.text === 'string' && Boolean(payload.text.trim());
}
if (endpoint === DELIVERY_ENDPOINTS.listTargets
|| endpoint === DELIVERY_ENDPOINTS.listSuggestions) {
return exactKeys(payload, ['botId']) && validBotId(payload.botId);
}
if (endpoint === DELIVERY_ENDPOINTS.createTarget) {
return exactKeys(payload, ['botId', 'target'])
&& validBotId(payload.botId) && validTarget(payload.target, { targetId: true });
}
if (endpoint === DELIVERY_ENDPOINTS.updateTarget) {
return exactKeys(payload, ['botId', 'targetId', 'target'])
&& validBotId(payload.botId) && validTargetId(payload.targetId)
&& validTarget(payload.target, { targetId: false });
}
if (endpoint === DELIVERY_ENDPOINTS.testTarget) {
return (exactKeys(payload, ['botId', 'targetId'])
&& validBotId(payload.botId) && validTargetId(payload.targetId))
|| (exactKeys(payload, ['botId', 'target'])
&& validBotId(payload.botId) && validDraftTarget(payload.target));
}
return exactKeys(payload, ['botId', 'targetId'])
&& validBotId(payload.botId) && validTargetId(payload.targetId);
}
function publicError(error) {
const code = PUBLIC_ERRORS.has(error?.code) ? error.code : 'delivery-failed';
return { code, message: code, details: {} };
}
export function createDeliveryRpcHandler(service) {
for (const method of [
'send',
'listTargets',
'listSuggestions',
'createTarget',
'updateTarget',
'deleteTarget',
]) {
if (typeof service?.[method] !== 'function') {
throw new TypeError(`A complete delivery service is required (${method})`);
}
}
return async (endpoint, payload, signal) => {
if (!validPayload(endpoint, payload)) {
return {
ok: false,
error: { code: 'bad-request', message: 'Invalid delivery request.', details: {} },
};
}
if (signal?.aborted) {
return { ok: false, error: { code: 'cancelled', message: 'cancelled', details: {} } };
}
try {
let value;
if (endpoint === DELIVERY_ENDPOINTS.send) {
value = await service.send(payload.botId, payload.targetId, payload.text, { signal });
} else if (endpoint === DELIVERY_ENDPOINTS.listTargets) {
value = await service.listTargets(payload.botId);
} else if (endpoint === DELIVERY_ENDPOINTS.listSuggestions) {
value = await service.listSuggestions(payload.botId);
} else if (endpoint === DELIVERY_ENDPOINTS.createTarget) {
value = await service.createTarget(payload.botId, payload.target);
} else if (endpoint === DELIVERY_ENDPOINTS.updateTarget) {
value = await service.updateTarget(payload.botId, payload.targetId, payload.target);
} else if (endpoint === DELIVERY_ENDPOINTS.deleteTarget) {
value = await service.deleteTarget(payload.botId, payload.targetId);
} else {
value = await service.send(
payload.botId,
Object.hasOwn(payload, 'target') ? payload.target : payload.targetId,
DELIVERY_TEST_MESSAGE,
{ signal },
);
}
return { ok: true, value };
} catch (error) {
return { ok: false, error: publicError(error) };
}
};
}
export function installDeliveryRpc(ctx, service, { authority } = {}) {
if (!ctx?.connection?.rpc || typeof ctx.connection.rpc.handle !== 'function') {
throw new TypeError('DSH Host Connection RPC is required');
}
return ctx.connection.rpc.handle(
DELIVERY_RPC_CHANNEL,
createDeliveryRpcHandler(service),
{ authority: resolveRpcAuthority(authority) },
);
}

View file

@ -0,0 +1,224 @@
import { normalizeDeliveryTarget } from './delivery-adapter.mjs';
const BOT_ID_PATTERN = /^[A-Za-z0-9_-]{1,128}$/;
const TARGET_ID_PATTERN = /^[A-Za-z0-9._:@-]{1,128}$/;
const CHANNEL_PATTERN = /^[a-z][a-z0-9-]{0,31}$/;
const DRAFT_TARGET_ID = '__test__';
const ADAPTER_METHODS = Object.freeze([
'ownsBot',
'listTargets',
'listSuggestions',
'createTarget',
'updateTarget',
'deleteTarget',
'sendText',
]);
const DELIVERY_ERROR_CODES = new Set([
'bad-request',
'unknown-bot',
'unknown-target',
'target-conflict',
'invalid-target',
'bot-not-connected',
'target-rejected',
'delivery-failed',
'cancelled',
]);
function deliveryError(code, message = code, options) {
const error = new Error(message, options);
error.code = code;
return error;
}
function botIdOf(value) {
if (typeof value !== 'string' || !BOT_ID_PATTERN.test(value)) {
throw deliveryError('bad-request', 'Invalid bot id');
}
return value;
}
function targetIdOf(value) {
if (typeof value !== 'string' || !TARGET_ID_PATTERN.test(value)) {
throw deliveryError('bad-request', 'Invalid target id');
}
return value;
}
function targetObject(value, { includesTargetId } = {}) {
if (!value || typeof value !== 'object' || Array.isArray(value)) {
throw deliveryError('bad-request', 'Invalid target');
}
if (includesTargetId) targetIdOf(value.targetId);
else if (Object.hasOwn(value, 'targetId')) {
throw deliveryError('bad-request', 'A target id cannot be changed');
}
return value;
}
function draftTargetObject(value) {
targetObject(value);
const keys = Object.keys(value);
if (keys.length !== 2 || !keys.includes('kind') || !keys.includes('route')) {
throw deliveryError('bad-request', 'Invalid draft target');
}
return value;
}
function cancellation(signal) {
if (signal?.aborted) throw deliveryError('cancelled', 'Request cancelled');
}
function publicOperationError(error, fallback = 'delivery-failed') {
if (error?.code === 'workspace-bot-not-found') {
return deliveryError('unknown-bot', 'Unknown bot', { cause: error });
}
if (DELIVERY_ERROR_CODES.has(error?.code)) return error;
return deliveryError(fallback, fallback, { cause: error });
}
function validateAdapter(adapter) {
if (!adapter || typeof adapter !== 'object'
|| typeof adapter.channel !== 'string' || !CHANNEL_PATTERN.test(adapter.channel)) {
throw new TypeError('A delivery adapter with a valid channel is required');
}
for (const method of ADAPTER_METHODS) {
if (typeof adapter[method] !== 'function') {
throw new TypeError(`A complete delivery adapter is required (${method})`);
}
}
return adapter;
}
export class DeliveryService {
#adapters = new Map();
registerAdapter(value) {
const adapter = validateAdapter(value);
const registration = Object.freeze({ adapter });
this.#adapters.set(adapter.channel, registration);
return () => {
if (this.#adapters.get(adapter.channel) !== registration) return false;
this.#adapters.delete(adapter.channel);
return true;
};
}
async listTargets(botId) {
const id = botIdOf(botId);
const adapter = await this.#adapterFor(id);
try {
const targets = await adapter.listTargets(id);
if (!Array.isArray(targets)) throw new TypeError('Adapter returned invalid targets');
return { botId: id, channel: adapter.channel, targets };
} catch (error) {
throw publicOperationError(error);
}
}
async listSuggestions(botId) {
const id = botIdOf(botId);
const adapter = await this.#adapterFor(id);
try {
const suggestions = await adapter.listSuggestions(id);
if (!Array.isArray(suggestions)) throw new TypeError('Adapter returned invalid suggestions');
return {
botId: id,
channel: adapter.channel,
suggestions: suggestions.map((suggestion) => normalizeDeliveryTarget(
adapter.channel,
suggestion,
{ targetIdRequired: false },
)),
};
} catch (error) {
throw publicOperationError(error);
}
}
async createTarget(botId, target) {
const id = botIdOf(botId);
targetObject(target, { includesTargetId: true });
const adapter = await this.#adapterFor(id);
try {
return await adapter.createTarget(id, target);
} catch (error) {
throw publicOperationError(error);
}
}
async updateTarget(botId, targetId, replacement) {
const id = botIdOf(botId);
const targetKey = targetIdOf(targetId);
targetObject(replacement);
const adapter = await this.#adapterFor(id);
try {
return await adapter.updateTarget(id, targetKey, replacement);
} catch (error) {
throw publicOperationError(error);
}
}
async deleteTarget(botId, targetId) {
const id = botIdOf(botId);
const targetKey = targetIdOf(targetId);
const adapter = await this.#adapterFor(id);
try {
await adapter.deleteTarget(id, targetKey);
return { deleted: true };
} catch (error) {
throw publicOperationError(error);
}
}
async send(botId, targetIdOrDraft, text, { signal } = {}) {
const id = botIdOf(botId);
const targetKey = typeof targetIdOrDraft === 'string'
? targetIdOf(targetIdOrDraft)
: null;
const draft = targetKey === null ? draftTargetObject(targetIdOrDraft) : null;
if (typeof text !== 'string' || !text.trim()) {
throw deliveryError('bad-request', 'Message text is required');
}
cancellation(signal);
const adapter = await this.#adapterFor(id);
try {
let target;
if (draft) {
target = { targetId: DRAFT_TARGET_ID, ...draft };
} else {
const targets = await adapter.listTargets(id);
if (!Array.isArray(targets)) throw new TypeError('Adapter returned invalid targets');
target = targets.find((candidate) => candidate?.targetId === targetKey);
if (!target) throw deliveryError('unknown-target', 'Unknown target');
}
cancellation(signal);
await adapter.sendText(id, target, text, { signal });
return { sent: true };
} catch (error) {
if (signal?.aborted || error?.name === 'AbortError' || error?.code === 'ABORT_ERR') {
throw deliveryError('cancelled', 'Request cancelled', { cause: error });
}
throw publicOperationError(error);
}
}
async #adapterFor(botId) {
for (const { adapter } of this.#adapters.values()) {
let ownsBot;
try {
ownsBot = await adapter.ownsBot(botId);
} catch (error) {
throw publicOperationError(error);
}
if (ownsBot) return adapter;
}
throw deliveryError('unknown-bot', 'Unknown bot');
}
}
export function createDeliveryService() {
return new DeliveryService();
}

View file

@ -0,0 +1,135 @@
const CHANNELS = new Set([
'weixin',
'feishu',
'dingtalk',
'wecom',
'qq',
'slack',
'telegram',
'discord',
'whatsapp',
]);
function isRecord(value) {
return value !== null && typeof value === 'object' && !Array.isArray(value);
}
function opaqueId(value) {
return typeof value === 'string'
&& value.length > 0
&& value.length <= 512
&& value.trim() === value
&& !/[\s:\u0000-\u001f\u007f]/u.test(value);
}
function afterPrefix(key, prefix) {
const marker = `${prefix}:`;
if (!key.startsWith(marker)) return null;
const value = key.slice(marker.length);
return opaqueId(value) ? value : null;
}
function simpleSuggestion(key, definitions) {
for (const [prefix, kind, field] of definitions) {
const value = afterPrefix(key, prefix);
if (value) return { kind, route: { [field]: value } };
}
return null;
}
function feishuSuggestion(key) {
const openId = afterPrefix(key, 'p2p');
if (openId?.startsWith('ou_')) {
return { kind: 'user', route: { openId } };
}
return simpleSuggestion(key, [['group', 'group', 'chatId']]);
}
function slackSuggestion(key) {
const directChannel = afterPrefix(key, 'direct');
if (directChannel && /^[A-Za-z0-9_-]{1,128}$/.test(directChannel)) {
return { kind: 'conversation', route: { channelId: directChannel } };
}
const match = /^group:([^:]+):(\d{1,20}(?:\.\d{1,20})?)$/.exec(key);
if (!match || !/^[A-Za-z0-9_-]{1,128}$/.test(match[1])) return null;
return { kind: 'thread', route: { channelId: match[1], threadTs: match[2] } };
}
function telegramSuggestion(key) {
const match = /^(direct|group):(-?\d+)(?::([1-9]\d*))?$/.exec(key);
if (!match) return null;
const chatId = Number(match[2]);
if (!Number.isSafeInteger(chatId)) return null;
if (match[1] === 'direct' && match[3] !== undefined) return null;
if (match[3] === undefined) return { kind: 'chat', route: { chatId: match[2] } };
const messageThreadId = Number(match[3]);
if (!Number.isSafeInteger(messageThreadId) || messageThreadId <= 0) return null;
return { kind: 'topic', route: { chatId: match[2], messageThreadId } };
}
function discordSuggestion(key) {
const match = /^(?:direct|group):(\d{1,32})$/.exec(key);
return match ? { kind: 'channel', route: { channelId: match[1] } } : null;
}
function whatsappSuggestion(key) {
const direct = /^direct:(\d{5,32}@(s\.whatsapp\.net|lid))$/.exec(key);
if (direct) return { kind: 'user', route: { jid: direct[1] } };
const group = /^group:(\d{5,32}(?:-\d{1,32})?@g\.us)$/.exec(key);
return group ? { kind: 'group', route: { jid: group[1] } } : null;
}
/** Convert one persisted conversation key into a stable proactive-delivery route. */
export function deliverySuggestionFromConversationKey(channel, key) {
if (!CHANNELS.has(channel) || typeof key !== 'string') return null;
switch (channel) {
case 'weixin':
return simpleSuggestion(key, [['p2p', 'user', 'toUserId']]);
case 'feishu':
return feishuSuggestion(key);
case 'dingtalk':
return simpleSuggestion(key, [
['p2p', 'user', 'userId'],
['group', 'group', 'openConversationId'],
]);
case 'wecom':
return simpleSuggestion(key, [
['direct', 'user', 'chatId'],
['group', 'group', 'chatId'],
]);
case 'qq':
return simpleSuggestion(key, [
['c2c', 'user', 'userOpenId'],
['group', 'group', 'groupOpenId'],
]);
case 'slack':
return slackSuggestion(key);
case 'telegram':
return telegramSuggestion(key);
case 'discord':
return discordSuggestion(key);
case 'whatsapp':
return whatsappSuggestion(key);
default:
return null;
}
}
/**
* Extract and de-duplicate stable delivery routes from a persisted sessions map.
* Session ids and every other state field are intentionally ignored.
*/
export function deliverySuggestionsFromSessions(channel, sessions) {
if (!CHANNELS.has(channel) || !isRecord(sessions)) return [];
const suggestions = [];
const seen = new Set();
for (const key of Object.keys(sessions)) {
const suggestion = deliverySuggestionFromConversationKey(channel, key);
if (!suggestion) continue;
const identity = JSON.stringify(suggestion);
if (seen.has(identity)) continue;
seen.add(identity);
suggestions.push(suggestion);
}
return suggestions;
}

View file

@ -10,6 +10,8 @@ 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 { createDeliveryService } from './delivery-service.mjs';
import { installUpdateRpc } from './update-rpc.mjs';
export const name = 'dsh-im-host';
@ -19,15 +21,18 @@ export const inject = [
'typertGateway',
];
function channelConfig(config, name) {
function channelConfig(config, name, deliveryService) {
const channel = config[name] ?? {};
return config.rpcAuthority === undefined
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 makeDeliveryService = internals.createDeliveryService ?? createDeliveryService;
const startFeishu = internals.applyFeishu ?? applyFeishu;
const startWeixin = internals.applyWeixin ?? applyWeixin;
const startDingtalk = internals.applyDingtalk ?? applyDingtalk;
@ -54,8 +59,17 @@ export function createImHostPlugin(internals = {}) {
name,
inject,
async apply(ctx, config = {}) {
const deliveryService = makeDeliveryService();
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,
}));
}
const activate = async (readyCtx) => {
await activateChannels(readyCtx, config);
await activateChannels(readyCtx, config, deliveryService);
};
if (typeof ctx?.inject === 'function') {
const modern = typeof ctx?.typertGateway?.stream === 'function';
@ -69,7 +83,7 @@ export function createImHostPlugin(internals = {}) {
},
});
async function activateChannels(ctx, config) {
async function activateChannels(ctx, config, deliveryService) {
setImHostLanguage(config.language ?? process.env.DSH_IM_LANGUAGE);
if (typeof ctx?.inject === 'function') {
ctx.inject(['tools', 'systemPrompt'], (artifactCtx) => {
@ -87,11 +101,16 @@ export function createImHostPlugin(internals = {}) {
} 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));
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);