mirror of
https://github.com/hansjone/dsh-im-ops.git
synced 2026-10-09 01:53:21 +08:00
Merge pull request #101 from evanfang0054/feat/feishu-topic-session-isolation
feat(feishu): isolate Harness sessions per topic in topic groups
This commit is contained in:
commit
7a16718922
5 changed files with 605 additions and 257 deletions
513
lib/index.js
513
lib/index.js
File diff suppressed because one or more lines are too long
|
|
@ -686,7 +686,7 @@ export class FeishuHarnessBridge {
|
|||
? pending.queue
|
||||
: null,
|
||||
isQuestionPending: () => this.#pendingInteractions.has(key),
|
||||
send: (text) => this.#send(event.message.chat_id, text),
|
||||
send: (text) => this.#send(event.message.chat_id, text, { replyTo: event.message.message_id }),
|
||||
});
|
||||
if (approvalReply) {
|
||||
const processing = approvalReply.process(async () => {
|
||||
|
|
@ -851,7 +851,7 @@ export class FeishuHarnessBridge {
|
|||
if (error?.code === 'turn-stopped') {
|
||||
await this.#removeProcessingReaction(messageId, processingReaction);
|
||||
if (error?.batchInputMessage) {
|
||||
await this.#send(event.message.chat_id, error.batchInputMessage).catch(() => undefined);
|
||||
await this.#send(event.message.chat_id, error.batchInputMessage, { replyTo: event.message.message_id }).catch(() => undefined);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
|
@ -940,7 +940,7 @@ export class FeishuHarnessBridge {
|
|||
]);
|
||||
}
|
||||
for (const reply of result?.messages ?? [result?.message]) {
|
||||
if (reply) await this.#send(event.message.chat_id, reply);
|
||||
if (reply) await this.#send(event.message.chat_id, reply, { replyTo: event.message.message_id });
|
||||
}
|
||||
this.#status.lastError = null;
|
||||
}
|
||||
|
|
@ -964,7 +964,7 @@ export class FeishuHarnessBridge {
|
|||
// accept() 侧已用 nonEmptyString(content) 判定,两侧保持一致。
|
||||
const commandText = !hasImages && !hasFiles && text ? text.trim() : null;
|
||||
if (!text && !hasImages && !hasFiles) {
|
||||
await this.#send(event.message.chat_id, t('目前支持文字、图片和文件消息。'));
|
||||
await this.#send(event.message.chat_id, t('目前支持文字、图片和文件消息。'), { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
|
||||
|
|
@ -973,11 +973,11 @@ export class FeishuHarnessBridge {
|
|||
return;
|
||||
}
|
||||
if (commandText === '/help') {
|
||||
await this.#send(event.message.chat_id, menuHelpText());
|
||||
await this.#send(event.message.chat_id, menuHelpText(), { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
if (MENU_COMMAND.test(commandText)) {
|
||||
await this.#sendMenuCard(key, event.message.chat_id);
|
||||
await this.#sendMenuCard(key, event.message.chat_id, { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
if (commandText === '/new') {
|
||||
|
|
@ -985,53 +985,54 @@ export class FeishuHarnessBridge {
|
|||
await this.#send(
|
||||
event.message.chat_id,
|
||||
t('当前任务仍在运行,请先停止任务或等待任务完成后再开启新会话。'),
|
||||
{ replyTo: event.message.message_id },
|
||||
);
|
||||
return;
|
||||
}
|
||||
await this.#state.clearSession(key);
|
||||
await this.#send(event.message.chat_id, t('已开启全新 Harness 会话。'));
|
||||
await this.#sendMenuCard(key, event.message.chat_id);
|
||||
await this.#send(event.message.chat_id, t('已开启全新 Harness 会话。'), { replyTo: event.message.message_id });
|
||||
await this.#sendMenuCard(key, event.message.chat_id, { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
if (commandText === '/status') {
|
||||
await this.#showStatusText(key, event.message.chat_id);
|
||||
await this.#showStatusText(key, event.message.chat_id, event.message.message_id);
|
||||
return;
|
||||
}
|
||||
if (commandText === '/compact') {
|
||||
const compactCommand = await runCompactCommand(commandText, this.#harness, this.#state, key, { signal: this.#signal });
|
||||
if (compactCommand) {
|
||||
await this.#send(event.message.chat_id, compactCommand.message);
|
||||
await this.#send(event.message.chat_id, compactCommand.message, { replyTo: event.message.message_id });
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (SESSION_LIST_PREFIX.test(commandText)) {
|
||||
const selector = commandText.replace(/^\/(?:sessionlist|sessions)/i, '').trim() || null;
|
||||
await this.#showSessions({ chatId: event.message.chat_id, key }, selector, 0);
|
||||
await this.#showSessions({ chatId: event.message.chat_id, key, replyTo: event.message.message_id }, selector, 0);
|
||||
return;
|
||||
}
|
||||
if (WORKSPACE_LIST_COMMAND.test(commandText)) {
|
||||
await this.#showWorkspaces({ chatId: event.message.chat_id, key });
|
||||
await this.#showWorkspaces({ chatId: event.message.chat_id, key, replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
if (WATCH_COMMAND.test(commandText)) {
|
||||
const target = (WATCH_COMMAND.exec(commandText)?.[1] ?? '').trim() || null;
|
||||
await this.#runWatch(key, event.message.chat_id, target);
|
||||
await this.#runWatch(key, event.message.chat_id, target, { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
if (UNWATCH_COMMAND.test(commandText)) {
|
||||
const target = (UNWATCH_COMMAND.exec(commandText)?.[1] ?? '').trim() || null;
|
||||
await this.#runUnwatch(key, event.message.chat_id, target);
|
||||
await this.#runUnwatch(key, event.message.chat_id, target, { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
if (WATCHLIST_COMMAND.test(commandText)) {
|
||||
await this.#showWatchList(key, event.message.chat_id);
|
||||
await this.#showWatchList(key, event.message.chat_id, { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
if (ARCHIVED_COMMAND.test(commandText)) {
|
||||
const match = ARCHIVED_COMMAND.exec(commandText);
|
||||
const value = match[1]?.toLowerCase();
|
||||
if (value !== 'on' && value !== 'off') {
|
||||
await this.#send(event.message.chat_id, t('用法:/archived on(包含归档会话)或 /archived off(隐藏归档会话)'));
|
||||
await this.#send(event.message.chat_id, t('用法:/archived on(包含归档会话)或 /archived off(隐藏归档会话)'), { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
if (typeof this.#state?.setIncludeArchivedSessions === 'function') {
|
||||
|
|
@ -1040,6 +1041,7 @@ export class FeishuHarnessBridge {
|
|||
await this.#send(
|
||||
event.message.chat_id,
|
||||
value === 'on' ? t('已开启:会话列表包含归档会话。') : t('已关闭:会话列表隐藏归档会话。'),
|
||||
{ replyTo: event.message.message_id },
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
|
@ -1059,7 +1061,7 @@ export class FeishuHarnessBridge {
|
|||
: await runWorkspaceCommand(text, this.#harness, key);
|
||||
if (workspaceCommand) {
|
||||
for (const reply of workspaceCommand.messages ?? [workspaceCommand.message]) {
|
||||
await this.#send(event.message.chat_id, reply);
|
||||
await this.#send(event.message.chat_id, reply, { replyTo: event.message.message_id });
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
|
@ -1073,7 +1075,7 @@ export class FeishuHarnessBridge {
|
|||
{ signal: this.#signal },
|
||||
);
|
||||
if (compactCommand) {
|
||||
await this.#send(event.message.chat_id, compactCommand.message);
|
||||
await this.#send(event.message.chat_id, compactCommand.message, { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
|
||||
|
|
@ -2030,7 +2032,7 @@ export class FeishuHarnessBridge {
|
|||
}
|
||||
|
||||
async #showSessions(
|
||||
{ chatId, key },
|
||||
{ chatId, key, replyTo = null },
|
||||
selector,
|
||||
page = 0,
|
||||
{ updateMessageId = null } = {},
|
||||
|
|
@ -2039,14 +2041,14 @@ export class FeishuHarnessBridge {
|
|||
const signal = this.#cardDataSignal();
|
||||
const resolved = await resolveSessionListWorkspace(selector ?? '', this.#harness, { signal });
|
||||
if (resolved.error) {
|
||||
await this.#send(chatId, resolved.error);
|
||||
await this.#send(chatId, resolved.error, { replyTo });
|
||||
return;
|
||||
}
|
||||
const listed = await this.#harness.listWorkspaceSessions(resolved.workspace, { signal });
|
||||
const sessions = this.#visibleSessions(Array.isArray(listed?.sessions) ? listed.sessions : []);
|
||||
const workspace = listed?.workspace ?? resolved.workspace;
|
||||
if (sessions.length === 0) {
|
||||
await this.#send(chatId, t('工作区:{workspace}\n该工作区暂无会话。', { workspace }));
|
||||
await this.#send(chatId, t('工作区:{workspace}\n该工作区暂无会话。', { workspace }), { replyTo });
|
||||
return;
|
||||
}
|
||||
const pageCount = Math.ceil(sessions.length / MENU_PAGE_SIZE);
|
||||
|
|
@ -2065,6 +2067,7 @@ export class FeishuHarnessBridge {
|
|||
{
|
||||
key,
|
||||
updateMessageId,
|
||||
replyTo,
|
||||
// Keep the canonical selector result for later page callbacks. The
|
||||
// list response's workspace is display data and is not authoritative.
|
||||
sessionWorkspace: resolved.workspace,
|
||||
|
|
@ -2076,7 +2079,7 @@ export class FeishuHarnessBridge {
|
|||
}
|
||||
}
|
||||
|
||||
async #showWorkspaces({ chatId, key }, { updateMessageId = null } = {}) {
|
||||
async #showWorkspaces({ chatId, key, replyTo = null }, { updateMessageId = null } = {}) {
|
||||
try {
|
||||
const { current, paths } = await workspacePathSnapshot(
|
||||
this.#harness,
|
||||
|
|
@ -2086,7 +2089,7 @@ export class FeishuHarnessBridge {
|
|||
await this.#sendCard(
|
||||
chatId,
|
||||
workspaceListCard(paths, current),
|
||||
{ key, updateMessageId },
|
||||
{ key, updateMessageId, replyTo },
|
||||
);
|
||||
} catch (error) {
|
||||
await this.#sendFailure(chatId, error, { logLabel: 'workspace list' });
|
||||
|
|
@ -2144,6 +2147,7 @@ export class FeishuHarnessBridge {
|
|||
|
||||
async #sendCard(chatId, cardJson, options = {}) {
|
||||
const updateMessageId = nonEmptyString(options.updateMessageId);
|
||||
const replyTo = nonEmptyString(options.replyTo);
|
||||
|
||||
if (updateMessageId) {
|
||||
try {
|
||||
|
|
@ -2161,6 +2165,28 @@ export class FeishuHarnessBridge {
|
|||
}
|
||||
}
|
||||
|
||||
// A brand-new card triggered by an inbound message is delivered as a
|
||||
// threaded reply so it lands inside the same Feishu topic. Falls back to
|
||||
// a plain chat message when the referenced message is gone.
|
||||
const content = cardJson;
|
||||
if (replyTo) {
|
||||
try {
|
||||
const response = await this.#client.im.v1.message.reply({
|
||||
path: { message_id: replyTo },
|
||||
data: { msg_type: 'interactive', content },
|
||||
});
|
||||
if (response?.code && response.code !== 0) {
|
||||
throw new Error(`Feishu card reply failed: ${response.msg || response.code}`);
|
||||
}
|
||||
const repliedMessageId = nonEmptyString(response?.data?.message_id);
|
||||
if (repliedMessageId) {
|
||||
this.#rememberCardRoute(repliedMessageId, chatId, options);
|
||||
return repliedMessageId;
|
||||
}
|
||||
} catch (error) {
|
||||
this.#logger.warn?.('[dsh-feishu] threaded card reply failed; sending a plain card:', error?.message ?? String(error));
|
||||
}
|
||||
}
|
||||
const response = await this.#client.im.v1.message.create({
|
||||
params: { receive_id_type: 'chat_id' },
|
||||
data: { receive_id: chatId, msg_type: 'interactive', content: cardJson },
|
||||
|
|
@ -2173,7 +2199,7 @@ export class FeishuHarnessBridge {
|
|||
return messageId;
|
||||
}
|
||||
|
||||
async #sendMenuCard(key, chatId, { updateMessageId = null } = {}) {
|
||||
async #sendMenuCard(key, chatId, { updateMessageId = null, replyTo = null } = {}) {
|
||||
let currentSessionId = null;
|
||||
let directSessionTitle = null;
|
||||
try {
|
||||
|
|
@ -2268,7 +2294,7 @@ export class FeishuHarnessBridge {
|
|||
currentSession: currentSessionId ? { id: currentSessionId, title: currentSessionTitle } : null,
|
||||
sessions, archiveVisible, presetCatalog, modelCatalog,
|
||||
}),
|
||||
{ key, updateMessageId },
|
||||
{ key, updateMessageId, replyTo },
|
||||
);
|
||||
}
|
||||
|
||||
|
|
@ -2350,7 +2376,7 @@ export class FeishuHarnessBridge {
|
|||
/**
|
||||
* Gather system status and show the status card.
|
||||
*/
|
||||
async #showStatusText(key, chatId) {
|
||||
async #showStatusText(key, chatId, replyTo = null) {
|
||||
try {
|
||||
await this.#harness.ensureRunning({ signal: this.#signal });
|
||||
const lines = [t('连接正常')];
|
||||
|
|
@ -2369,7 +2395,7 @@ export class FeishuHarnessBridge {
|
|||
: (settings.agentPreset || t('跟随默认')),
|
||||
}));
|
||||
}
|
||||
await this.#send(chatId, lines.join('\n'));
|
||||
await this.#send(chatId, lines.join('\n'), { replyTo });
|
||||
} catch (error) {
|
||||
await this.#sendFailure(chatId, error, { logLabel: 'status text' });
|
||||
}
|
||||
|
|
@ -2787,11 +2813,12 @@ export class FeishuHarnessBridge {
|
|||
notify = true,
|
||||
validatedTarget = null,
|
||||
workspaceHint = null,
|
||||
replyTo = null,
|
||||
} = {}) {
|
||||
const watchRequestedAt = Date.now();
|
||||
const reply = async (message) => {
|
||||
if (!notify) return;
|
||||
await this.#send(chatId, message).catch((error) => {
|
||||
await this.#send(chatId, message, { replyTo }).catch((error) => {
|
||||
this.#logger.warn?.('[dsh-feishu] watch notification failed:', error.message);
|
||||
});
|
||||
};
|
||||
|
|
@ -2854,10 +2881,10 @@ export class FeishuHarnessBridge {
|
|||
return { ok: true, changed: !existingEntry, entry: this.#state.watchEntry?.(key, resolved.sessionId) };
|
||||
}
|
||||
|
||||
async #runUnwatch(key, chatId, target, { notify = true } = {}) {
|
||||
async #runUnwatch(key, chatId, target, { notify = true, replyTo = null } = {}) {
|
||||
const reply = async (message) => {
|
||||
if (!notify) return;
|
||||
await this.#send(chatId, message).catch((error) => {
|
||||
await this.#send(chatId, message, { replyTo }).catch((error) => {
|
||||
this.#logger.warn?.('[dsh-feishu] unwatch notification failed:', error.message);
|
||||
});
|
||||
};
|
||||
|
|
@ -2883,7 +2910,7 @@ export class FeishuHarnessBridge {
|
|||
return { ok: true, changed: true, entry };
|
||||
}
|
||||
|
||||
async #showWatchList(key, chatId, { updateMessageId = null } = {}) {
|
||||
async #showWatchList(key, chatId, { updateMessageId = null, replyTo = null } = {}) {
|
||||
const entries = this.#state.watchEntries?.(key) ?? [];
|
||||
// 收集可选会话(用于「添加关注」多选下拉);失败则传空数组 → 只渲染移除/列表。
|
||||
let availableSessions = [];
|
||||
|
|
@ -2913,6 +2940,7 @@ export class FeishuHarnessBridge {
|
|||
{
|
||||
key,
|
||||
updateMessageId,
|
||||
replyTo,
|
||||
sessionWorkspace: currentWorkspace,
|
||||
},
|
||||
);
|
||||
|
|
@ -3068,17 +3096,18 @@ export class FeishuHarnessBridge {
|
|||
actor: senderOpenId(event),
|
||||
chatId: event.message.chat_id,
|
||||
requiresMention: event.message.chat_type !== 'p2p',
|
||||
replyToMessageId: event.message.message_id,
|
||||
}),
|
||||
onInteractionResolved: (resolution) => this.#handleInteractionResolved(resolution),
|
||||
files,
|
||||
};
|
||||
}
|
||||
|
||||
async #sendAnswerText(chatId, answer, { deliveryId, presentation }) {
|
||||
async #sendAnswerText(chatId, answer, { deliveryId, presentation, replyTo = null }) {
|
||||
const providerMessageIds = [];
|
||||
for (const chunk of splitText(answer)) {
|
||||
this.#signal?.throwIfAborted();
|
||||
const messageId = await this.#send(chatId, chunk);
|
||||
const messageId = await this.#send(chatId, chunk, { replyTo });
|
||||
if (messageId) providerMessageIds.push(messageId);
|
||||
}
|
||||
return createDeliveryReceipt({
|
||||
|
|
@ -3182,6 +3211,7 @@ export class FeishuHarnessBridge {
|
|||
{
|
||||
deliveryId: messageId,
|
||||
presentation: 'feishu-text',
|
||||
replyTo: messageId,
|
||||
},
|
||||
);
|
||||
} catch (error) {
|
||||
|
|
@ -3252,6 +3282,7 @@ export class FeishuHarnessBridge {
|
|||
{
|
||||
deliveryId: messageId,
|
||||
presentation: 'feishu-text-fallback',
|
||||
replyTo: messageId,
|
||||
},
|
||||
);
|
||||
} catch (fallbackError) {
|
||||
|
|
@ -3302,6 +3333,7 @@ export class FeishuHarnessBridge {
|
|||
{
|
||||
deliveryId: messageId,
|
||||
presentation: 'feishu-text-fallback',
|
||||
replyTo: messageId,
|
||||
},
|
||||
);
|
||||
} catch (fallbackError) {
|
||||
|
|
@ -3368,11 +3400,11 @@ export class FeishuHarnessBridge {
|
|||
const pending = this.#pendingInteractions.get(key);
|
||||
if (!pending || pending !== expected || pending.submitting) {
|
||||
if (this.#isResolvedQuestionReply(event, key)) {
|
||||
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT()).catch(() => undefined);
|
||||
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT(), { replyTo: event.message.message_id }).catch(() => undefined);
|
||||
return;
|
||||
}
|
||||
if (claimed && (!pending || pending !== expected)) {
|
||||
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT());
|
||||
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT(), { replyTo: event.message.message_id });
|
||||
return;
|
||||
}
|
||||
return this.#enqueueMessage(event, messageId, key, processingReaction, {
|
||||
|
|
@ -3430,7 +3462,7 @@ export class FeishuHarnessBridge {
|
|||
if (error?.code === 'interaction-not-pending') {
|
||||
this.#rememberResolvedInteraction(key, pending);
|
||||
this.#clearPendingInteraction(key, pending.interactionId);
|
||||
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT()).catch(() => undefined);
|
||||
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT(), { replyTo: event.message.message_id }).catch(() => undefined);
|
||||
return;
|
||||
}
|
||||
pending.submitting = false;
|
||||
|
|
@ -3448,12 +3480,13 @@ export class FeishuHarnessBridge {
|
|||
actor,
|
||||
chatId,
|
||||
requiresMention,
|
||||
replyToMessageId,
|
||||
}) {
|
||||
if (await this.#approvals.handleRequested(interaction, {
|
||||
key,
|
||||
actor,
|
||||
requiresMention,
|
||||
send: (text) => this.#send(chatId, text),
|
||||
send: (text) => this.#send(chatId, text, { replyTo: replyToMessageId }),
|
||||
})) return;
|
||||
|
||||
// Approval requests return above; the existing question state machine stays unchanged.
|
||||
|
|
@ -3484,6 +3517,7 @@ export class FeishuHarnessBridge {
|
|||
await this.#send(
|
||||
chatId,
|
||||
t('检测到这个 Session 中遗留的待回答问题,已安全取消并继续处理你刚才的消息。'),
|
||||
{ replyTo: replyToMessageId },
|
||||
).catch(() => undefined);
|
||||
return;
|
||||
}
|
||||
|
|
@ -3519,6 +3553,7 @@ export class FeishuHarnessBridge {
|
|||
answers: [],
|
||||
index: 0,
|
||||
chatId,
|
||||
replyToMessageId,
|
||||
queue: null,
|
||||
claimedReplyMessageId: null,
|
||||
submitting: false,
|
||||
|
|
@ -3553,6 +3588,9 @@ export class FeishuHarnessBridge {
|
|||
pending.questions.length,
|
||||
{ requiresMention: pending.requiresMention },
|
||||
),
|
||||
// Reply to the message that started the turn so the question lands in
|
||||
// the same Feishu thread/topic instead of the group's default area.
|
||||
{ replyTo: pending.replyToMessageId },
|
||||
);
|
||||
if (messageId) {
|
||||
pending.questionMessageIds.add(messageId);
|
||||
|
|
@ -3585,7 +3623,7 @@ export class FeishuHarnessBridge {
|
|||
await this.#state.markSeen(messageId);
|
||||
this.#status.lastMessageAt = new Date().toISOString();
|
||||
this.#status.messagesReceived += 1;
|
||||
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT()).catch(() => undefined);
|
||||
await this.#send(event.message.chat_id, INTERACTION_RESOLVED_TEXT(), { replyTo: event.message.message_id }).catch(() => undefined);
|
||||
}
|
||||
|
||||
#takePendingInteraction(key, interactionId) {
|
||||
|
|
@ -3651,13 +3689,33 @@ export class FeishuHarnessBridge {
|
|||
else processingReaction.success();
|
||||
}
|
||||
|
||||
async #send(chatId, text) {
|
||||
async #send(chatId, text, { replyTo } = {}) {
|
||||
const content = JSON.stringify({ text });
|
||||
if (replyTo) {
|
||||
try {
|
||||
const response = await this.#client.im.v1.message.reply({
|
||||
path: { message_id: replyTo },
|
||||
data: { msg_type: 'text', content },
|
||||
});
|
||||
if (response?.code && response.code !== 0) {
|
||||
throw new Error(`Feishu reply failed: ${response.msg || response.code}`);
|
||||
}
|
||||
return nonEmptyString(response?.data?.message_id);
|
||||
} catch (error) {
|
||||
// The referenced message may be gone (recalled/deleted); keep the
|
||||
// delivery promise by falling back to a plain chat message.
|
||||
this.#logger.warn?.(
|
||||
'[dsh-feishu] threaded reply failed; falling back to a plain message:',
|
||||
error?.message ?? String(error),
|
||||
);
|
||||
}
|
||||
}
|
||||
const response = await this.#client.im.v1.message.create({
|
||||
params: { receive_id_type: 'chat_id' },
|
||||
data: {
|
||||
receive_id: chatId,
|
||||
msg_type: 'text',
|
||||
content: JSON.stringify({ text }),
|
||||
content,
|
||||
},
|
||||
});
|
||||
if (response?.code && response.code !== 0) {
|
||||
|
|
|
|||
|
|
@ -16,6 +16,11 @@ export function conversationKey(event) {
|
|||
}
|
||||
const chatId = event?.message?.chat_id;
|
||||
if (!chatId) throw new Error('Feishu group event has no chat id');
|
||||
// Topic groups: every message belongs to a thread, so key the session per
|
||||
// thread to keep each topic's Harness conversation isolated. Regular group
|
||||
// chats carry no thread_id and keep the single shared `group:<chat_id>` key.
|
||||
const threadId = event?.message?.thread_id;
|
||||
if (typeof threadId === 'string' && threadId.trim()) return `group:${chatId}:thread:${threadId}`;
|
||||
return `group:${chatId}`;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -1238,6 +1238,180 @@ test('a threaded Feishu reply answers a pending Harness question before the orig
|
|||
assert.equal(status.messagesReplied, 1);
|
||||
});
|
||||
|
||||
test('a Harness question is presented as a threaded reply inside a topic group', async () => {
|
||||
const sent = [];
|
||||
const replied = [];
|
||||
const streamed = [];
|
||||
const seen = new Set();
|
||||
const sessions = new Map();
|
||||
const status = {
|
||||
messagesReceived: 0,
|
||||
messagesReplied: 0,
|
||||
messagesRejected: 0,
|
||||
lastMessageAt: null,
|
||||
lastReplyAt: null,
|
||||
lastRejectedAt: null,
|
||||
lastError: null,
|
||||
};
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: {
|
||||
im: { v1: { message: {
|
||||
create: async (request) => {
|
||||
sent.push({ text: JSON.parse(request.data.content).text });
|
||||
return { code: 0, data: { message_id: `om_sent_${sent.length}` } };
|
||||
},
|
||||
reply: async (request) => {
|
||||
replied.push({
|
||||
to: request.path.message_id,
|
||||
text: JSON.parse(request.data.content).text,
|
||||
});
|
||||
return { code: 0, data: { message_id: `om_replied_${replied.length}` } };
|
||||
},
|
||||
} } },
|
||||
},
|
||||
channel: {
|
||||
stream: async (_chatId, input) => {
|
||||
await input.markdown({
|
||||
setContent: async (content) => streamed.push(content),
|
||||
});
|
||||
return { messageId: 'om_stream' };
|
||||
},
|
||||
},
|
||||
harness: {
|
||||
ensureRunning: async () => true,
|
||||
sessionExists: async () => false,
|
||||
createSession: async () => 'session-topic',
|
||||
ask: async (_sessionId, _text, options) => {
|
||||
await options.onUpdate({ type: 'tool', name: 'ask_user_question' });
|
||||
await options.onInteraction({
|
||||
kind: 'question',
|
||||
interactionId: 'question-rpc',
|
||||
rpcId: 'question-rpc',
|
||||
sessionId: 'session-topic',
|
||||
payload: {
|
||||
type: 'question/requested',
|
||||
sessionId: 'session-topic',
|
||||
questions: [{
|
||||
id: 'environment',
|
||||
header: '测试环境',
|
||||
question: '请选择测试环境',
|
||||
options: [{ label: '测试环境' }, { label: '生产环境' }],
|
||||
}],
|
||||
},
|
||||
respond: async () => ({ accepted: true }),
|
||||
});
|
||||
return '已完成';
|
||||
},
|
||||
},
|
||||
state: {
|
||||
hasSeen: (id) => seen.has(id),
|
||||
markSeen: async (id) => seen.add(id),
|
||||
sessionFor: (key) => sessions.get(key) ?? null,
|
||||
setSession: async (key, sessionId) => sessions.set(key, sessionId),
|
||||
clearSession: async (key) => sessions.delete(key),
|
||||
},
|
||||
status,
|
||||
allowedSenderOpenIds: new Set(['ou_user']),
|
||||
});
|
||||
|
||||
bridge.accept(event('om_prompt', '请先调用 ask_user_question', {
|
||||
chat_type: 'group',
|
||||
chat_id: 'oc_topic_group',
|
||||
thread_id: 'omt_prompt',
|
||||
}));
|
||||
await bridge.waitForIdle();
|
||||
|
||||
assert.deepEqual(
|
||||
replied,
|
||||
[{ to: 'om_prompt', text: replied[0]?.text }],
|
||||
'the question must be delivered through the reply API targeting the triggering message',
|
||||
);
|
||||
assert.ok(replied[0]?.text.includes('请选择测试环境'));
|
||||
assert.equal(
|
||||
sent.some(({ text }) => text.includes('请选择测试环境')),
|
||||
false,
|
||||
'the question must not be sent as a plain chat message that lands outside the topic',
|
||||
);
|
||||
assert.equal(streamed.at(-1), '已完成');
|
||||
});
|
||||
|
||||
test('command replies in a topic group are threaded to the triggering message', async () => {
|
||||
const created = [];
|
||||
const replied = [];
|
||||
const seen = new Set();
|
||||
const sessions = new Map([['group:oc_topic:thread:omt_cmd', 'session-cmd']]);
|
||||
const bridge = new FeishuHarnessBridge({
|
||||
client: {
|
||||
im: { v1: { message: {
|
||||
create: async (request) => {
|
||||
created.push({
|
||||
type: request.data.msg_type,
|
||||
text: request.data.msg_type === 'text'
|
||||
? JSON.parse(request.data.content).text
|
||||
: null,
|
||||
});
|
||||
return { code: 0, data: { message_id: `om_created_${created.length}` } };
|
||||
},
|
||||
reply: async (request) => {
|
||||
replied.push({
|
||||
to: request.path.message_id,
|
||||
type: request.data.msg_type,
|
||||
text: request.data.msg_type === 'text'
|
||||
? JSON.parse(request.data.content).text
|
||||
: null,
|
||||
});
|
||||
return { code: 0, data: { message_id: `om_replied_${replied.length}` } };
|
||||
},
|
||||
} } },
|
||||
},
|
||||
harness: {
|
||||
ensureRunning: async () => true,
|
||||
sessionExists: async () => true,
|
||||
},
|
||||
state: {
|
||||
hasSeen: (id) => seen.has(id),
|
||||
markSeen: async (id) => seen.add(id),
|
||||
sessionFor: (key) => sessions.get(key) ?? null,
|
||||
setSession: async (key, sessionId) => sessions.set(key, sessionId),
|
||||
clearSession: async (key) => sessions.delete(key),
|
||||
},
|
||||
status: {
|
||||
messagesReceived: 0,
|
||||
messagesReplied: 0,
|
||||
messagesRejected: 0,
|
||||
lastMessageAt: null,
|
||||
lastReplyAt: null,
|
||||
lastRejectedAt: null,
|
||||
lastError: null,
|
||||
},
|
||||
allowedSenderOpenIds: new Set(['ou_user']),
|
||||
});
|
||||
|
||||
bridge.accept(event('om_cmd', '/new', {
|
||||
chat_type: 'group',
|
||||
chat_id: 'oc_topic',
|
||||
thread_id: 'omt_cmd',
|
||||
}));
|
||||
await eventually(
|
||||
() => replied.length >= 2,
|
||||
'the /new replies were not presented as threaded replies',
|
||||
);
|
||||
|
||||
assert.ok(
|
||||
replied.some(({ to, text }) => to === 'om_cmd' && text?.includes('已开启全新')),
|
||||
'the /new confirmation must be delivered through the reply API targeting the command message',
|
||||
);
|
||||
assert.ok(
|
||||
replied.some(({ to, type }) => to === 'om_cmd' && type === 'interactive'),
|
||||
'the menu card must also be threaded to the command message',
|
||||
);
|
||||
assert.equal(
|
||||
created.some(({ text }) => text?.includes('已开启全新')),
|
||||
false,
|
||||
'the /new confirmation must not be sent as a plain chat message outside the topic',
|
||||
);
|
||||
});
|
||||
|
||||
test('pending Harness questions are isolated by Feishu conversation', async () => {
|
||||
const fixture = stateFixture([
|
||||
['p2p:ou_a', 'session-a'],
|
||||
|
|
|
|||
|
|
@ -232,6 +232,38 @@ test('conversationKey isolates p2p users and groups', () => {
|
|||
}), 'group:oc_group');
|
||||
});
|
||||
|
||||
test('conversationKey isolates topic-group threads without affecting regular groups', () => {
|
||||
// Topic groups: each message belongs to a thread, so every topic gets its own session.
|
||||
assert.equal(conversationKey({
|
||||
sender: { sender_id: { open_id: 'ou_test' } },
|
||||
message: { chat_type: 'group', chat_id: 'oc_topic_group', thread_id: 'om_thread_a' },
|
||||
}), 'group:oc_topic_group:thread:om_thread_a');
|
||||
assert.equal(conversationKey({
|
||||
sender: { sender_id: { open_id: 'ou_other' } },
|
||||
message: { chat_type: 'group', chat_id: 'oc_topic_group', thread_id: 'om_thread_b' },
|
||||
}), 'group:oc_topic_group:thread:om_thread_b');
|
||||
assert.notEqual(
|
||||
conversationKey({
|
||||
sender: { sender_id: { open_id: 'ou_test' } },
|
||||
message: { chat_type: 'group', chat_id: 'oc_topic_group', thread_id: 'om_thread_a' },
|
||||
}),
|
||||
conversationKey({
|
||||
sender: { sender_id: { open_id: 'ou_test' } },
|
||||
message: { chat_type: 'group', chat_id: 'oc_topic_group', thread_id: 'om_thread_b' },
|
||||
}),
|
||||
);
|
||||
// Regular group chats: no thread_id keeps the single shared group key.
|
||||
assert.equal(conversationKey({
|
||||
sender: { sender_id: { open_id: 'ou_test' } },
|
||||
message: { chat_type: 'group', chat_id: 'oc_group' },
|
||||
}), 'group:oc_group');
|
||||
// Blank thread_id values fall back to the shared group key.
|
||||
assert.equal(conversationKey({
|
||||
sender: { sender_id: { open_id: 'ou_test' } },
|
||||
message: { chat_type: 'group', chat_id: 'oc_group', thread_id: ' ' },
|
||||
}), 'group:oc_group');
|
||||
});
|
||||
|
||||
test('splitText preserves all text', () => {
|
||||
const input = `${'a'.repeat(12)}\n${'b'.repeat(12)}`;
|
||||
const chunks = splitText(input, 15);
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue