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 <cursoragent@cursor.com>
This commit is contained in:
oliver 2026-09-06 12:45:36 +08:00
parent 356538ac18
commit 01924e5f7e
8 changed files with 307 additions and 0 deletions

View file

@ -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();

View file

@ -170,6 +170,7 @@ export async function createProductionController(ctx, config = {}, internals = {
}).start();
return {
controller,
stateFor,
deliveryAdapter: createDeliveryAdapter({
channel: 'whatsapp', workspaces, coreController, stateFor,
}),

View file

@ -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:<jid>: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<string|null>} 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) => {