From 01924e5f7eba4d1d4a6e7e5bea01b693e6252a98 Mon Sep 17 00:00:00 2001 From: oliver Date: Sun, 6 Sep 2026 12:45:36 +0800 Subject: [PATCH] Expose conversation/target session lookup for cron result mirroring. Schedulers can resolve the live Harness session for a bot+target or conversationKey without changing proactive send() semantics. Co-authored-by: Cursor --- PROACTIVE_DELIVERY.en.md | 2 + PROACTIVE_DELIVERY.md | 2 + plugin-src/host/channels/whatsapp/index.mjs | 32 +++++ .../host/channels/whatsapp/production.mjs | 1 + plugin-src/host/index.mjs | 89 +++++++++++++ .../shared/conversation-state-store.mjs | 13 ++ src/channels/shared/delivery-session-keys.mjs | 42 ++++++ test/delivery-session-resolve.test.mjs | 126 ++++++++++++++++++ 8 files changed, 307 insertions(+) create mode 100644 src/channels/shared/delivery-session-keys.mjs create mode 100644 test/delivery-session-resolve.test.mjs diff --git a/PROACTIVE_DELIVERY.en.md b/PROACTIVE_DELIVERY.en.md index 8639f1a..be2f2ad 100644 --- a/PROACTIVE_DELIVERY.en.md +++ b/PROACTIVE_DELIVERY.en.md @@ -316,6 +316,8 @@ No. You can edit its name, type, and native route without changing call paramete A `sessionId` identifies a Harness Session. It is not a uniform, stable message address across the nine platforms. Proactive delivery uses the stable bot and saved-target pair instead. +Same-Host schedulers that need the delivery summary in a follow-up-capable Harness session should call `resolveConversationSession` / `resolveTargetSession` (channel-registered binding lookups) and then `session.append` / `inject` themselves. `send()` does not write the session log. + ### The test succeeded, but the recipient cannot see the message A successful test proves only that the platform accepted the send. Check bot permissions, platform restrictions, target accuracy, and client-side filtering or archive settings. diff --git a/PROACTIVE_DELIVERY.md b/PROACTIVE_DELIVERY.md index 75d73aa..00182f3 100644 --- a/PROACTIVE_DELIVERY.md +++ b/PROACTIVE_DELIVERY.md @@ -316,6 +316,8 @@ Connection RPC 默认只允许当前 Host 的回环调用。若 Web profile 明 `sessionId` 标识 Harness 会话,不是九个平台统一、稳定的消息投递地址。主动投递只使用机器人和已保存目标的稳定组合。 +同 Host 的调度插件若需要把投递摘要写回可继续聊的 Harness 会话,应另行调用 `resolveConversationSession` / `resolveTargetSession`(由渠道注册的会话绑定反查),再自行 `session.append` / `inject`;`send()` 本身不写会话日志。 + ### 测试成功但对方没有看到消息 测试成功只证明平台接口接受发送。请继续检查机器人权限、平台限制、目标是否正确,以及客户端侧的消息过滤或归档设置。 diff --git a/plugin-src/host/channels/whatsapp/index.mjs b/plugin-src/host/channels/whatsapp/index.mjs index ecebe94..419c7da 100644 --- a/plugin-src/host/channels/whatsapp/index.mjs +++ b/plugin-src/host/channels/whatsapp/index.mjs @@ -1,5 +1,6 @@ import { createProductionController } from './production.mjs'; import { installWhatsappRpc } from './rpc.mjs'; +import { conversationKeyMatchesTarget } from '../../../../src/channels/shared/delivery-session-keys.mjs'; export const name = 'dsh-im-whatsapp-host'; export const inject = ['connection', 'typertGateway']; @@ -18,6 +19,7 @@ export async function apply(ctx, config = {}) { config.rpcAuthority, ); let unregisterPeer = () => {}; + let unregisterConversation = () => {}; try { const dshIm = typeof ctx.get === 'function' ? ctx.get('dshIm') : undefined; if (dshIm && typeof dshIm.registerPeerResolver === 'function' @@ -26,10 +28,40 @@ export async function apply(ctx, config = {}) { production.controller.resolveChannelPeer(sessionId) )); } + if (dshIm && typeof dshIm.registerConversationSessionResolver === 'function' + && typeof production.stateFor === 'function') { + const resolve = async (botId, conversationKey) => { + try { + const state = await production.stateFor(botId); + const sessionId = state?.sessionFor?.(conversationKey); + return typeof sessionId === 'string' && sessionId ? sessionId : null; + } catch { + return null; + } + }; + resolve.findByTarget = async (botId, target) => { + try { + const state = await production.stateFor(botId); + const rows = typeof state?.listLiveSessions === 'function' + ? state.listLiveSessions() + : []; + for (const row of rows) { + if (conversationKeyMatchesTarget(row.conversationKey, target)) { + return row; + } + } + } catch { + return null; + } + return null; + }; + unregisterConversation = dshIm.registerConversationSessionResolver(resolve); + } } catch { // dshIm optional during partial boots } ctx.effect(() => async () => { + unregisterConversation?.(); unregisterPeer?.(); await unregisterDelivery?.(); await production.close(); diff --git a/plugin-src/host/channels/whatsapp/production.mjs b/plugin-src/host/channels/whatsapp/production.mjs index 58a7ff8..355c006 100644 --- a/plugin-src/host/channels/whatsapp/production.mjs +++ b/plugin-src/host/channels/whatsapp/production.mjs @@ -170,6 +170,7 @@ export async function createProductionController(ctx, config = {}, internals = { }).start(); return { controller, + stateFor, deliveryAdapter: createDeliveryAdapter({ channel: 'whatsapp', workspaces, coreController, stateFor, }), diff --git a/plugin-src/host/index.mjs b/plugin-src/host/index.mjs index bdd65b5..6927bff 100644 --- a/plugin-src/host/index.mjs +++ b/plugin-src/host/index.mjs @@ -10,6 +10,9 @@ 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 { + conversationKeysFromDeliveryTarget, +} from '../../src/channels/shared/delivery-session-keys.mjs'; import { installDeliveryRpc } from './delivery-rpc.mjs'; import { installDeliveryHttp } from './delivery-http.mjs'; import { createDeliveryService } from './delivery-service.mjs'; @@ -63,6 +66,68 @@ export function createImHostPlugin(internals = {}) { async apply(ctx, config = {}) { const deliveryService = makeDeliveryService(); const peerResolvers = []; + const conversationSessionResolvers = []; + + async function resolveConversationSession(botId, conversationKey) { + const id = typeof botId === 'string' ? botId.trim() : ''; + const key = typeof conversationKey === 'string' ? conversationKey.trim() : ''; + if (!id || !key) return null; + for (const resolve of conversationSessionResolvers) { + try { + const sessionId = await resolve(id, key); + if (typeof sessionId === 'string' && sessionId.trim()) { + return { sessionId: sessionId.trim(), botId: id, conversationKey: key }; + } + } catch { + // try next channel + } + } + return null; + } + + /** + * Prefer an exact conversationKey binding; for group targets also accept + * the first live `group::user:…` binding when scope is user_in_chat. + */ + async function resolveTargetSession(botId, targetId) { + const id = typeof botId === 'string' ? botId.trim() : ''; + const tid = typeof targetId === 'string' ? targetId.trim() : ''; + if (!id || !tid) return null; + let targets = []; + try { + const listed = await deliveryService.listTargets(id); + targets = Array.isArray(listed?.targets) ? listed.targets : []; + } catch { + return null; + } + const target = targets.find((row) => row?.targetId === tid); + if (!target) return null; + + const keys = conversationKeysFromDeliveryTarget(target); + for (const key of keys) { + const hit = await resolveConversationSession(id, key); + if (hit) return hit; + } + + // Group deliveries may only have per-speaker bindings (user_in_chat). + for (const resolve of conversationSessionResolvers) { + if (typeof resolve.findByTarget !== 'function') continue; + try { + const hit = await resolve.findByTarget(id, target); + const sessionId = typeof hit?.sessionId === 'string' ? hit.sessionId.trim() : ''; + const conversationKey = typeof hit?.conversationKey === 'string' + ? hit.conversationKey.trim() + : ''; + if (sessionId && conversationKey) { + return { sessionId, botId: id, conversationKey }; + } + } catch { + // try next + } + } + return null; + } + if (typeof ctx?.provide === 'function') { ctx.provide('dshIm', Object.freeze({ send: (botId, targetId, text, options) => ( @@ -97,6 +162,30 @@ export function createImHostPlugin(internals = {}) { if (index >= 0) peerResolvers.splice(index, 1); }; }, + /** + * Resolve the live Harness session bound to a conversationKey for a bot. + * @param {string} botId + * @param {string} conversationKey + */ + resolveConversationSession, + /** + * Resolve the live Harness session for a proactive-delivery target. + * @param {string} botId + * @param {string} targetId + */ + resolveTargetSession, + /** + * @param {(botId: string, conversationKey: string) => Promise} resolve + * Optional `resolve.findByTarget(botId, target)` for group prefix matches. + */ + registerConversationSessionResolver: (resolve) => { + if (typeof resolve !== 'function') return () => {}; + conversationSessionResolvers.push(resolve); + return () => { + const index = conversationSessionResolvers.indexOf(resolve); + if (index >= 0) conversationSessionResolvers.splice(index, 1); + }; + }, })); } const activate = async (readyCtx) => { diff --git a/src/channels/shared/conversation-state-store.mjs b/src/channels/shared/conversation-state-store.mjs index f3e80e7..332d751 100644 --- a/src/channels/shared/conversation-state-store.mjs +++ b/src/channels/shared/conversation-state-store.mjs @@ -100,6 +100,19 @@ export class ConversationStateStore { return this.#state.sessions[key] ?? null; } + /** + * Live conversationKey → sessionId bindings (excludes unbound provenance). + * @returns {Array<{ conversationKey: string, sessionId: string }>} + */ + listLiveSessions() { + return Object.entries(this.#state.sessions) + .filter(([conversationKey, sessionId]) => ( + typeof conversationKey === 'string' && conversationKey + && typeof sessionId === 'string' && sessionId + )) + .map(([conversationKey, sessionId]) => ({ conversationKey, sessionId })); + } + /** * Reverse-map a Harness session id to its conversation binding key. * Live bindings win; unbound sessions still resolve via sessionOrigins diff --git a/src/channels/shared/delivery-session-keys.mjs b/src/channels/shared/delivery-session-keys.mjs new file mode 100644 index 0000000..1d7f50e --- /dev/null +++ b/src/channels/shared/delivery-session-keys.mjs @@ -0,0 +1,42 @@ +/** + * Map a proactive-delivery target route to conversationKey candidates + * so schedulers can resolve the live Harness session for that chat. + */ + +/** + * @param {object|null|undefined} target + * @returns {string[]} + */ +export function conversationKeysFromDeliveryTarget(target) { + if (!target || typeof target !== 'object') return []; + const kind = typeof target.kind === 'string' ? target.kind.trim() : ''; + const route = target.route && typeof target.route === 'object' ? target.route : {}; + const jid = String(route.jid || route.chatId || route.openId || route.userId || '') + .trim(); + if (!jid) return []; + + if (kind === 'group' || jid.endsWith('@g.us')) { + const groupJid = jid.toLowerCase().endsWith('@g.us') ? jid : `${jid}`; + return [`group:${groupJid}`]; + } + + // DM / user targets: prefer the native jid; phone-shaped ids already include @s.whatsapp.net. + return [`direct:${jid}`]; +} + +/** + * Whether a live conversationKey belongs to a delivery target. + * Groups accept both shared (`group:jid`) and per-user (`group:jid:user:…`) bindings. + * @param {string} conversationKey + * @param {object|null|undefined} target + */ +export function conversationKeyMatchesTarget(conversationKey, target) { + const key = typeof conversationKey === 'string' ? conversationKey.trim() : ''; + if (!key) return false; + const candidates = conversationKeysFromDeliveryTarget(target); + for (const candidate of candidates) { + if (key === candidate) return true; + if (candidate.startsWith('group:') && key.startsWith(`${candidate}:user:`)) return true; + } + return false; +} diff --git a/test/delivery-session-resolve.test.mjs b/test/delivery-session-resolve.test.mjs new file mode 100644 index 0000000..6e2da01 --- /dev/null +++ b/test/delivery-session-resolve.test.mjs @@ -0,0 +1,126 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; + +import { + conversationKeyMatchesTarget, + conversationKeysFromDeliveryTarget, +} from '../src/channels/shared/delivery-session-keys.mjs'; +import { createImHostPlugin } from '../plugin-src/host/index.mjs'; + +const CHANNELS = [ + ['feishu', 'applyFeishu'], + ['weixin', 'applyWeixin'], + ['dingtalk', 'applyDingtalk'], + ['wecom', 'applyWecom'], + ['qq', 'applyQq'], + ['slack', 'applySlack'], + ['telegram', 'applyTelegram'], + ['discord', 'applyDiscord'], + ['whatsapp', 'applyWhatsapp'], + ['office', 'applyOffice'], +]; + +test('conversationKeysFromDeliveryTarget maps WhatsApp DM and group routes', () => { + assert.deepEqual( + conversationKeysFromDeliveryTarget({ + kind: 'user', + route: { jid: '8613800000000@s.whatsapp.net' }, + }), + ['direct:8613800000000@s.whatsapp.net'], + ); + assert.deepEqual( + conversationKeysFromDeliveryTarget({ + kind: 'group', + route: { jid: '120363@g.us' }, + }), + ['group:120363@g.us'], + ); +}); + +test('conversationKeyMatchesTarget accepts shared and per-user group keys', () => { + const target = { kind: 'group', route: { jid: '120363@g.us' } }; + assert.equal(conversationKeyMatchesTarget('group:120363@g.us', target), true); + assert.equal( + conversationKeyMatchesTarget('group:120363@g.us:user:86138@s.whatsapp.net', target), + true, + ); + assert.equal(conversationKeyMatchesTarget('direct:86138@s.whatsapp.net', target), false); +}); + +test('dshIm resolveConversationSession and resolveTargetSession use registered lookups', async () => { + const targets = [ + { targetId: 'wa-dm', kind: 'user', route: { jid: '86138@s.whatsapp.net' } }, + { targetId: 'ops-group', kind: 'group', route: { jid: '120363@g.us' } }, + ]; + const deliveryService = { + async send() { return { sent: true }; }, + async listTargets(botId) { + return { botId, channel: 'whatsapp', targets }; + }, + async listCatalog() { return { bots: [] }; }, + async createTarget() { return {}; }, + }; + const provided = []; + const internals = Object.fromEntries(CHANNELS.map(([, applyName]) => [ + applyName, + async () => {}, + ])); + Object.assign(internals, { + createDeliveryService: () => deliveryService, + installUpdateRpc: () => {}, + installDeliveryRpc: () => {}, + installDeliveryHttp: () => {}, + }); + const ctx = { + connection: { rpc: {} }, + webServer: { register() {} }, + effect() {}, + provide: (...args) => provided.push(args), + }; + + await createImHostPlugin(internals).apply(ctx, { rpcAuthority: 'trusted-host' }); + const dshIm = provided[0][1]; + + const resolve = async (botId, conversationKey) => { + if (botId === 'bot_1' && conversationKey === 'direct:86138@s.whatsapp.net') { + return 'sess-dm'; + } + return null; + }; + resolve.findByTarget = async (botId, target) => { + if (botId === 'bot_1' && target?.targetId === 'ops-group') { + return { + conversationKey: 'group:120363@g.us:user:alice@lid', + sessionId: 'sess-group-user', + }; + } + return null; + }; + dshIm.registerConversationSessionResolver(resolve); + + assert.deepEqual( + await dshIm.resolveConversationSession('bot_1', 'direct:86138@s.whatsapp.net'), + { + sessionId: 'sess-dm', + botId: 'bot_1', + conversationKey: 'direct:86138@s.whatsapp.net', + }, + ); + assert.deepEqual( + await dshIm.resolveTargetSession('bot_1', 'wa-dm'), + { + sessionId: 'sess-dm', + botId: 'bot_1', + conversationKey: 'direct:86138@s.whatsapp.net', + }, + ); + assert.deepEqual( + await dshIm.resolveTargetSession('bot_1', 'ops-group'), + { + sessionId: 'sess-group-user', + botId: 'bot_1', + conversationKey: 'group:120363@g.us:user:alice@lid', + }, + ); + assert.equal(await dshIm.resolveTargetSession('bot_1', 'missing'), null); +});