mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-09 00:33:20 +08:00
Plugins inject conversation.session.header.actions; resolve sessionId back to the bound chat and render nickname/phone or group title. Co-authored-by: Cursor <cursoragent@cursor.com>
293 lines
13 KiB
JavaScript
293 lines
13 KiB
JavaScript
import QRCode from 'qrcode';
|
|
|
|
import { publicConnectionTestResult } from '../../../../src/channels/shared/connection-test.mjs';
|
|
import { SET_ACCESS_POLICY_ENDPOINT, validAccessPolicyPayload } from '../shared/access-policy-rpc.mjs';
|
|
import {
|
|
REFRESH_ACCESS_GROUP_TITLES_ENDPOINT,
|
|
RESOLVE_ACCESS_PENDING_ENDPOINT,
|
|
SET_ACCESS_GRANT_ENDPOINT,
|
|
validAccessGrantPayload,
|
|
validAccessGroupTitlesRefreshPayload,
|
|
validAccessPendingResolvePayload,
|
|
} from '../shared/access-grant-rpc.mjs';
|
|
import { SET_GROUP_SESSION_SCOPE_ENDPOINT, validGroupSessionScopePayload } from '../shared/group-session-scope-rpc.mjs';
|
|
import { SET_CONTEXT_ENHANCEMENT_ENDPOINT, validContextEnhancementPayload } from '../shared/context-enhancement-rpc.mjs';
|
|
import { resolveRpcAuthority } from '../../rpc-authority.mjs';
|
|
import { publicWorkspaceError, SET_WORKSPACE_ENDPOINT, validWorkspacePayload } from '../shared/workspace-rpc.mjs';
|
|
import { SET_AGENT_PRESET_ENDPOINT, validAgentPresetPayload } from '../shared/agent-preset-rpc.mjs';
|
|
|
|
export const WHATSAPP_RPC_CHANNEL = '/whatsapp';
|
|
export const WHATSAPP_ENDPOINTS = Object.freeze({
|
|
status: 'connection.status',
|
|
beginProvisioning: 'provision.begin',
|
|
pollProvisioning: 'provision.poll',
|
|
cancelProvisioning: 'provision.cancel',
|
|
reconnectBot: 'bot.reconnect',
|
|
deleteBot: 'bot.delete',
|
|
setAccessPolicy: SET_ACCESS_POLICY_ENDPOINT,
|
|
setAccessGrant: SET_ACCESS_GRANT_ENDPOINT,
|
|
resolveAccessPending: RESOLVE_ACCESS_PENDING_ENDPOINT,
|
|
refreshAccessGroupTitles: REFRESH_ACCESS_GROUP_TITLES_ENDPOINT,
|
|
setGroupSessionScope: SET_GROUP_SESSION_SCOPE_ENDPOINT,
|
|
setWorkspace: SET_WORKSPACE_ENDPOINT,
|
|
setAgentPreset: SET_AGENT_PRESET_ENDPOINT,
|
|
setContextEnhancement: SET_CONTEXT_ENHANCEMENT_ENDPOINT,
|
|
resolveChannelPeer: 'bot.session.channel-peer',
|
|
});
|
|
export const WHATSAPP_RPC_ENDPOINTS = Object.freeze(Object.values(WHATSAPP_ENDPOINTS));
|
|
|
|
const FORBIDDEN_PUBLIC_KEYS = new Set(['qrValue', 'accountJid', 'authDirectory']);
|
|
|
|
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 ([WHATSAPP_ENDPOINTS.status, WHATSAPP_ENDPOINTS.beginProvisioning].includes(endpoint)) {
|
|
return exactKeys(payload, []) ? null : `${endpoint} does not accept fields.`;
|
|
}
|
|
if ([WHATSAPP_ENDPOINTS.pollProvisioning, WHATSAPP_ENDPOINTS.cancelProvisioning].includes(endpoint)) {
|
|
return exactKeys(payload, ['attemptId']) && validId(payload.attemptId)
|
|
? null : `${endpoint} requires an attemptId.`;
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.reconnectBot) {
|
|
return exactKeys(payload, ['botId', 'sendTest']) && validId(payload.botId)
|
|
&& (payload.sendTest === undefined || typeof payload.sendTest === 'boolean')
|
|
? null : 'bot.reconnect requires a botId and optional sendTest flag.';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.deleteBot) {
|
|
return exactKeys(payload, ['botId', 'confirm']) && validId(payload.botId)
|
|
&& payload.confirm === true ? null : 'bot.delete requires a botId and confirm=true.';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.setAccessPolicy) {
|
|
return validAccessPolicyPayload(payload)
|
|
? null : '请提交有效的访问设置。';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.setAccessGrant) {
|
|
return validAccessGrantPayload(payload)
|
|
? null : '请提交有效的分级访问授权。';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.resolveAccessPending) {
|
|
return validAccessPendingResolvePayload(payload)
|
|
? null : '请提交有效的审批请求。';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.refreshAccessGroupTitles) {
|
|
return validAccessGroupTitlesRefreshPayload(payload)
|
|
? null : '请提交有效的群名同步请求。';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.setGroupSessionScope) {
|
|
return validGroupSessionScopePayload(payload)
|
|
? null : '请提交有效的群会话策略。';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.setWorkspace) {
|
|
return validWorkspacePayload(payload)
|
|
? null : '请输入工作区绝对路径。';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.setAgentPreset) {
|
|
return validAgentPresetPayload(payload)
|
|
? null : '请选择 Agent Preset。';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.setContextEnhancement) {
|
|
return validContextEnhancementPayload(payload)
|
|
? null : '请提交有效的上下文增强设置。';
|
|
}
|
|
if (endpoint === WHATSAPP_ENDPOINTS.resolveChannelPeer) {
|
|
return exactKeys(payload, ['sessionId'])
|
|
&& typeof payload.sessionId === 'string'
|
|
&& payload.sessionId.trim()
|
|
&& payload.sessionId.trim().length <= 256
|
|
? null : 'bot.session.channel-peer requires a sessionId.';
|
|
}
|
|
return 'Unknown WhatsApp endpoint.';
|
|
}
|
|
|
|
function sanitizePublic(value) {
|
|
if (Array.isArray(value)) return value.map(sanitizePublic);
|
|
if (!isRecord(value)) return value;
|
|
const safe = {};
|
|
for (const [key, child] of Object.entries(value)) {
|
|
if (!FORBIDDEN_PUBLIC_KEYS.has(key)) safe[key] = sanitizePublic(child);
|
|
}
|
|
return safe;
|
|
}
|
|
|
|
async function qrDataUrl(value) {
|
|
return QRCode.toDataURL(value, {
|
|
type: 'image/png',
|
|
errorCorrectionLevel: 'M',
|
|
margin: 2,
|
|
width: 320,
|
|
});
|
|
}
|
|
|
|
async function encodeAttempt(value, encodeQr) {
|
|
if (!value || typeof value.qrValue !== 'string') return sanitizePublic(value);
|
|
return sanitizePublic({ ...value, qrCodeDataUrl: await encodeQr(value.qrValue) });
|
|
}
|
|
|
|
async function publicStatus(value, encodeQr) {
|
|
const snapshot = structuredClone(value);
|
|
if (snapshot?.provisioning) snapshot.provisioning = await encodeAttempt(snapshot.provisioning, encodeQr);
|
|
return sanitizePublic(snapshot);
|
|
}
|
|
|
|
export function createWhatsappRpcHandler(controller, { encodeQr = qrDataUrl } = {}) {
|
|
for (const method of ['status', 'startProvisioning', 'registrationStatus', 'cancelProvisioning', 'reconnectBot', 'deleteBot']) {
|
|
if (typeof controller?.[method] !== 'function') {
|
|
throw new TypeError(`A complete WhatsApp controller is required (${method})`);
|
|
}
|
|
}
|
|
const qrCache = new Map();
|
|
const cachedEncode = (value) => {
|
|
let encoded = qrCache.get(value);
|
|
if (!encoded) {
|
|
if (qrCache.size >= 16) qrCache.delete(qrCache.keys().next().value);
|
|
encoded = Promise.resolve().then(() => encodeQr(value));
|
|
qrCache.set(value, encoded);
|
|
}
|
|
return encoded;
|
|
};
|
|
return async (endpoint, payload, signal) => {
|
|
if (signal?.aborted) return { ok: false, error: { code: 'cancelled', message: 'The request was cancelled.' } };
|
|
if (!WHATSAPP_RPC_ENDPOINTS.includes(endpoint)) {
|
|
return { ok: false, error: { code: 'bad-request', message: 'Unknown WhatsApp endpoint.' } };
|
|
}
|
|
const invalid = payloadFailure(endpoint, payload);
|
|
if (invalid) return { ok: false, error: { code: 'bad-request', message: invalid } };
|
|
try {
|
|
let value;
|
|
if (endpoint === WHATSAPP_ENDPOINTS.status) {
|
|
value = await publicStatus(await controller.status(), cachedEncode);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.beginProvisioning) {
|
|
value = await encodeAttempt(await controller.startProvisioning(), cachedEncode);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.pollProvisioning) {
|
|
const attempt = await controller.registrationStatus(payload.attemptId);
|
|
if (!attempt) return { ok: false, error: { code: 'bad-request', message: 'The provisioning attempt no longer exists.' } };
|
|
value = await encodeAttempt(attempt, cachedEncode);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.cancelProvisioning) {
|
|
value = sanitizePublic(await controller.cancelProvisioning(payload.attemptId));
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.reconnectBot) {
|
|
const checked = await controller.reconnectBot(payload.botId);
|
|
if (signal?.aborted) {
|
|
return { ok: false, error: { code: 'cancelled', message: 'The request was cancelled.' } };
|
|
}
|
|
value = await publicStatus(checked, cachedEncode);
|
|
if (payload.sendTest === true) {
|
|
let testError = null;
|
|
const connected = checked?.bots?.some(
|
|
(bot) => bot?.botId === payload.botId && bot.connected === true,
|
|
) === true;
|
|
if (!connected) {
|
|
testError = new Error('WhatsApp bot is not connected');
|
|
testError.code = 'test-target-unavailable';
|
|
} else {
|
|
try {
|
|
if (typeof controller.sendConnectionTest !== 'function') {
|
|
const unavailable = new Error('Connection test is unavailable');
|
|
unavailable.code = 'test-target-unavailable';
|
|
throw unavailable;
|
|
}
|
|
await controller.sendConnectionTest(payload.botId);
|
|
} catch (error) {
|
|
testError = error;
|
|
}
|
|
}
|
|
value = { ...value, testMessage: publicConnectionTestResult(testError) };
|
|
}
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.setWorkspace) {
|
|
if (typeof controller.updateWorkspace !== 'function') throw new Error('Workspace update is unavailable');
|
|
value = await publicStatus(
|
|
await controller.updateWorkspace(payload.botId, payload.workspace),
|
|
cachedEncode,
|
|
);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.setContextEnhancement) {
|
|
if (typeof controller.updateContextEnhancement !== 'function') throw new Error('Context enhancement update is unavailable');
|
|
value = await controller.updateContextEnhancement(
|
|
payload.botId, payload.config, (status) => publicStatus(status, cachedEncode),
|
|
);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.setAgentPreset) {
|
|
if (typeof controller.updateAgentPreset !== 'function') throw new Error('Agent preset update is unavailable');
|
|
value = await publicStatus(
|
|
await controller.updateAgentPreset(payload.botId, payload.agentPreset),
|
|
cachedEncode,
|
|
);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.setAccessPolicy) {
|
|
if (typeof controller.updateAccessPolicy !== 'function') throw new Error('Access policy update is unavailable');
|
|
value = await controller.updateAccessPolicy(
|
|
payload.botId,
|
|
payload.policy,
|
|
(status) => publicStatus(status, cachedEncode),
|
|
);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.setAccessGrant) {
|
|
if (typeof controller.updateAccessGrant !== 'function') throw new Error('Access grant update is unavailable');
|
|
value = await controller.updateAccessGrant(
|
|
payload.botId,
|
|
payload.grant,
|
|
(status) => publicStatus(status, cachedEncode),
|
|
);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.resolveAccessPending) {
|
|
if (typeof controller.resolveAccessPending !== 'function') throw new Error('Access pending resolve is unavailable');
|
|
value = await controller.resolveAccessPending(
|
|
payload.botId,
|
|
{
|
|
pendingId: payload.pendingId,
|
|
action: payload.action,
|
|
resolvedByPhone: payload.resolvedByPhone,
|
|
},
|
|
(status) => publicStatus(status, cachedEncode),
|
|
);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.refreshAccessGroupTitles) {
|
|
if (typeof controller.refreshAccessGroupTitles !== 'function') {
|
|
throw new Error('Group title refresh is unavailable');
|
|
}
|
|
value = await publicStatus(
|
|
await controller.refreshAccessGroupTitles(payload.botId, payload.groupJids),
|
|
cachedEncode,
|
|
);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.setGroupSessionScope) {
|
|
if (typeof controller.updateGroupSessionScope !== 'function') throw new Error('Group session scope update is unavailable');
|
|
value = await controller.updateGroupSessionScope(
|
|
payload.botId,
|
|
payload.groupSessionScope,
|
|
(status) => publicStatus(status, cachedEncode),
|
|
);
|
|
} else if (endpoint === WHATSAPP_ENDPOINTS.resolveChannelPeer) {
|
|
if (typeof controller.resolveChannelPeer !== 'function') {
|
|
throw new Error('Channel peer resolve is unavailable');
|
|
}
|
|
value = await controller.resolveChannelPeer(payload.sessionId.trim());
|
|
} else {
|
|
value = await publicStatus(await controller.deleteBot(payload.botId), cachedEncode);
|
|
}
|
|
return signal?.aborted
|
|
? { ok: false, error: { code: 'cancelled', message: 'The request was cancelled.' } }
|
|
: { ok: true, value };
|
|
} catch (error) {
|
|
const workspaceError = publicWorkspaceError(error);
|
|
return signal?.aborted
|
|
? { ok: false, error: { code: 'cancelled', message: 'The request was cancelled.' } }
|
|
: { ok: false, error: workspaceError
|
|
?? { code: 'whatsapp-operation-failed', message: 'WhatsApp 操作失败,请稍后重试。' } };
|
|
}
|
|
};
|
|
}
|
|
|
|
export function installWhatsappRpc(ctx, controller, options, 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(
|
|
WHATSAPP_RPC_CHANNEL,
|
|
createWhatsappRpcHandler(controller, options),
|
|
{ authority: resolveRpcAuthority(authority) },
|
|
);
|
|
}
|