diff --git a/src/bot-registry.ts b/src/bot-registry.ts index f2a1e0e6f..ec92e7a70 100644 --- a/src/bot-registry.ts +++ b/src/bot-registry.ts @@ -35,6 +35,40 @@ export type ChatReplyMode = 'chat' | 'new-topic' | 'shared' | 'chat-topic'; export type ContentTriggerScope = 'topic' | 'regularGroup' | 'both'; export type ContentTriggerMatchType = 'keyword' | 'regex'; export type ContentTriggerActionType = 'start-or-wake-session'; +export type MessageListenerSenderType = 'user' | 'bot'; + +export interface MessageListenerConfig { + enabled: boolean; + name?: string; + replyCardTitle?: string; + workingDir?: string; + prompt: string; + senderPolicy?: { + /** + * all_except_excluded: listen to all matching sender types except excluded ids. + * include_only: listen only to includeSenderOpenIds; empty include means none. + */ + mode?: 'all_except_excluded' | 'include_only'; + includeSenderOpenIds?: string[]; + excludeSenderOpenIds?: string[]; + includeSenderTypes?: MessageListenerSenderType[]; + excludeSenderTypes?: MessageListenerSenderType[]; + /** Default true. */ + excludeSelf?: boolean; + }; + messagePolicy?: { + /** Defaults to text + post. */ + includeMsgTypes?: string[]; + /** V1 only supports top-level group messages. */ + scope?: 'top_level'; + }; + replyPolicy?: { + /** V1 always replies under the triggering message. */ + mode?: 'thread'; + /** V1 starts one session per matched message. */ + sessionMode?: 'per_message'; + }; +} export interface SummaryRangeConfig { /** 0 means no count limit; omitted defaults to 50. */ @@ -94,6 +128,12 @@ function normalizeContentTriggerScope(raw: unknown): ContentTriggerScope | undef return undefined; } +function normalizeMessageListenerSenderType(raw: unknown): MessageListenerSenderType | undefined { + if (raw === 'user' || raw === 'human') return 'user'; + if (raw === 'bot' || raw === 'app') return 'bot'; + return undefined; +} + function normalizeNonNegativeInt(raw: unknown): number | undefined { if (typeof raw !== 'number') return undefined; if (!Number.isInteger(raw) || raw < 0) return undefined; @@ -750,6 +790,76 @@ function normalizeContentTriggers(raw: unknown, botIndex: number): ContentTrigge return out.length > 0 ? out : undefined; } +function normalizeMessageListenerStringList(raw: unknown): string[] | undefined { + const values = normalizeStringList(raw); + return values.length > 0 ? [...new Set(values)] : undefined; +} + +function normalizeMessageListenerSenderTypes(raw: unknown): MessageListenerSenderType[] | undefined { + if (!Array.isArray(raw)) return undefined; + const values = raw + .map(normalizeMessageListenerSenderType) + .filter((value): value is MessageListenerSenderType => !!value); + return values.length > 0 ? [...new Set(values)] : undefined; +} + +function normalizeMessageListenerConfig(raw: unknown, botIndex: number, chatId: string): MessageListenerConfig | undefined { + if (!raw || typeof raw !== 'object' || Array.isArray(raw)) return undefined; + const entry = raw as Record; + const prompt = normalizeNonEmptyString(entry.prompt); + const enabled = entry.enabled === true; + if (enabled && !prompt) { + logger.warn(`Bot config [${botIndex}] messageListeners[${chatId}] ignored: enabled listener requires prompt`); + return undefined; + } + + const senderRaw = entry.senderPolicy && typeof entry.senderPolicy === 'object' && !Array.isArray(entry.senderPolicy) + ? entry.senderPolicy as Record + : {}; + const senderPolicy: NonNullable = {}; + const mode = senderRaw.mode === 'include_only' ? 'include_only' : 'all_except_excluded'; + const includeSenderOpenIds = normalizeMessageListenerStringList(senderRaw.includeSenderOpenIds); + const excludeSenderOpenIds = normalizeMessageListenerStringList(senderRaw.excludeSenderOpenIds); + const includeSenderTypes = normalizeMessageListenerSenderTypes(senderRaw.includeSenderTypes); + const excludeSenderTypes = normalizeMessageListenerSenderTypes(senderRaw.excludeSenderTypes); + if (mode !== 'all_except_excluded') senderPolicy.mode = mode; + if (includeSenderOpenIds) senderPolicy.includeSenderOpenIds = includeSenderOpenIds; + if (excludeSenderOpenIds) senderPolicy.excludeSenderOpenIds = excludeSenderOpenIds; + if (includeSenderTypes) senderPolicy.includeSenderTypes = includeSenderTypes; + if (excludeSenderTypes) senderPolicy.excludeSenderTypes = excludeSenderTypes; + if (senderRaw.excludeSelf === false) senderPolicy.excludeSelf = false; + + const messageRaw = entry.messagePolicy && typeof entry.messagePolicy === 'object' && !Array.isArray(entry.messagePolicy) + ? entry.messagePolicy as Record + : {}; + const messagePolicy: NonNullable = {}; + const includeMsgTypes = normalizeMessageListenerStringList(messageRaw.includeMsgTypes); + if (includeMsgTypes) messagePolicy.includeMsgTypes = includeMsgTypes; + messagePolicy.scope = 'top_level'; + + return { + enabled, + ...(normalizeNonEmptyString(entry.name) ? { name: normalizeNonEmptyString(entry.name) } : {}), + ...(normalizeNonEmptyString(entry.replyCardTitle) ? { replyCardTitle: normalizeNonEmptyString(entry.replyCardTitle) } : {}), + ...(normalizeNonEmptyString(entry.workingDir) ? { workingDir: normalizeNonEmptyString(entry.workingDir) } : {}), + prompt: prompt ?? '', + ...(Object.keys(senderPolicy).length > 0 ? { senderPolicy } : {}), + ...(Object.keys(messagePolicy).length > 0 ? { messagePolicy } : {}), + replyPolicy: { mode: 'thread', sessionMode: 'per_message' }, + }; +} + +function normalizeMessageListeners(raw: unknown, botIndex: number): Record | undefined { + if (!raw || typeof raw !== 'object' || Array.isArray(raw)) return undefined; + const out: Record = {}; + for (const [chatId, listenerRaw] of Object.entries(raw as Record)) { + if (typeof chatId !== 'string' || !chatId.trim()) continue; + const listener = normalizeMessageListenerConfig(listenerRaw, botIndex, chatId.trim()); + if (listener) out[chatId.trim()] = listener; + } + return Object.keys(out).length > 0 ? out : undefined; +} + export interface OncallChat { /** Lark chat_id (oc_xxx) the bot was pulled into. */ chatId: string; @@ -1254,6 +1364,14 @@ export interface BotConfig { * Default (undefined) = passive. */ autoStartOnNewTopic?: boolean; + /** + * Per-chat group message listener. Keyed by chat_id and bot-scoped so the + * dashboard can configure it from the Roles page's natural group × bot + * matrix. When enabled, the bot may react to non-@ top-level group messages + * after deterministic sender/msgType filtering. V1 always replies in a + * fresh thread under the triggering message. + */ + messageListeners?: Record; /** * Worktree picker mode on the repo-select card. When true, the worktree * control renders the multi-repo selector (pick N repos + branch) instead of @@ -2051,6 +2169,7 @@ export function parseBotConfigsFromText(jsonText: string): BotConfig[] { : undefined; const summaryRange = normalizeSummaryRange(entry.summaryRange ?? entry.summary); const contentTriggers = normalizeContentTriggers(entry.contentTriggers, i); + const messageListeners = normalizeMessageListeners(entry.messageListeners, i); const vcMeetingAgent = normalizeVcMeetingAgentConfig(entry.vcMeetingAgent); // voice:per-bot 语音引擎覆盖。结构化保留(engine ∈ sami|openai,sami/openai @@ -2179,6 +2298,7 @@ export function parseBotConfigsFromText(jsonText: string): BotConfig[] { ? entry.autoStartOnGroupJoinPrompt : undefined, autoStartOnNewTopic: entry.autoStartOnNewTopic === true || undefined, + messageListeners, worktreeMultiPicker: entry.worktreeMultiPicker === true || undefined, // Per-bot regular-group default mode. Only the non-default modes // ('chat-topic' | 'new-topic' | 'shared') are meaningful; 'chat' (the flat diff --git a/src/core/dashboard-ipc-server.ts b/src/core/dashboard-ipc-server.ts index 2f22e5f4c..4cc384c09 100644 --- a/src/core/dashboard-ipc-server.ts +++ b/src/core/dashboard-ipc-server.ts @@ -68,8 +68,8 @@ import * as chatFirstSeenStore from '../services/chat-first-seen-store.js'; import * as scheduler from './scheduler.js'; import { listActiveSessions, findActiveBySessionId, closeSession, getActiveSessionsRegistry, transferSession, deliverWriteLinkCardToOwners, forkWorker, suspendWorker, killWorker } from './worker-pool.js'; import { listOnlineDaemons } from '../utils/daemon-discovery.js'; +import { getChatMode, replyMessage, sendMessage, resolveUnionIdFromOpenId, listThreadMessages, listChatMessages, listChatMessagesUntil, listChatBotMembers, getUserProfile, getUserProfileStrict, resolveAllowedUsersWithMap, type ChatBotMember } from '../im/lark/client.js'; import { isSessionStopped } from './session-liveness.js'; -import { getChatMode, replyMessage, sendMessage, resolveUnionIdFromOpenId, listThreadMessages, listChatMessages, listChatBotMembers, getUserProfile, getUserProfileStrict, resolveAllowedUsersWithMap, type ChatBotMember } from '../im/lark/client.js'; import { parseApiMessage, cardContentHasUpgradeFallback, resolveMergedCardContent } from '../im/lark/message-parser.js'; import { resumeSession, spawnDashboardSession, activateQueuedSession, closeCliMismatchedSessionsForBot, suspendActiveSessionsForBot } from './session-manager.js'; import { parseSpawnRequest } from './session-create.js'; @@ -95,6 +95,24 @@ import { validateTriggerRequest, type TriggerResponse } from '../services/trigge import { resolveCliSelection, selectionKeyForBot } from '../setup/cli-selection.js'; import { checkCliAvailability } from '../setup/cli-availability.js'; import { enrichHistorySenders, type HistoryBotInfo } from '../dashboard/history-senders.js'; +import { getMessageListenerConfig, sanitizeMessageListenerUpdate, updateMessageListenerConfig, validateMessageListenerUpdate } from '../services/message-listener-store.js'; +import { + MAX_MESSAGE_LISTENER_PROMPT_BYTES, + normalizeMessageListenerPreviewLimit, + previewMessageListenerMatches, + renderMessageListenerInstruction, + type MessageListenerPreviewMatch, +} from '../services/message-listener.js'; +import { + createMessageListenerRunPreview, + createMessageListenerRunPreviewTurnId, + getMessageListenerRunPreview, + markMessageListenerRunPreviewFailed, + markMessageListenerRunPreviewTriggered, +} from '../services/message-listener-run-preview-store.js'; +import { listChatMemberDisplays } from '../services/groups-store.js'; + +const MESSAGE_LISTENER_PREVIEW_WINDOW_MS = 24 * 60 * 60 * 1000; let exactChatGrantHandler: typeof applyExactChatGrantRequest = applyExactChatGrantRequest; /** Test seam: replace the exact-grant service without touching live Feishu/config state. */ @@ -129,7 +147,7 @@ import { getBotName, type SessionRow, } from './dashboard-rows.js'; -import { getBotBrand, getBot, loadBotConfigs, readBotSkillPolicy, getBotTuiSlashAllow } from '../bot-registry.js'; +import { getBotBrand, getBot, loadBotConfigs, readBotSkillPolicy, getBotTuiSlashAllow, type MessageListenerConfig } from '../bot-registry.js'; import { normalizeKanbanColumn, normalizeKanbanPosition, normalizeSessionTitle } from './session-board.js'; import { validateSlashInjection } from './slash-inject.js'; import { validateRoleLibraryPath } from './role-library.js'; @@ -1933,12 +1951,13 @@ ipcRoute('GET', '/api/groups', async (_req, res) => { const enriched = chats.map(c => { const oncall = oncallStore.getOncallStatus(cachedLarkAppId, c.chatId); const hasRole = resolveRoleFile(cachedLarkAppId, c.chatId) !== null; + const hasMessageListener = getMessageListenerConfig(cachedLarkAppId, c.chatId)?.enabled === true; // /introduce 记录的外部 botmux 机器人(按名字)——dashboard 团队看板用 // 它识别「介绍过同团队机器人的协作群」。 const observedBotNames = observedBotsStore .listObservedBots(config.session.dataDir, cachedLarkAppId, c.chatId) .map(b => b.name); - return { ...c, oncallChat: oncall ?? null, firstSeenAt: seenMap.get(c.chatId) ?? null, hasRole, observedBotNames }; + return { ...c, oncallChat: oncall ?? null, firstSeenAt: seenMap.get(c.chatId) ?? null, hasRole, hasMessageListener, observedBotNames }; }); jsonRes(res, 200, { chats: enriched }); } catch (e) { @@ -2148,6 +2167,264 @@ ipcRoute('DELETE', '/api/roles/:chatId', async (_req, res, p) => { jsonRes(res, 200, { ok: true, existed }); }); +ipcRoute('GET', '/api/message-listeners/:chatId', async (_req, res, p) => { + if (!cachedLarkAppId) return jsonRes(res, 503, { error: 'larkAppId_not_set' }); + if (!isValidRoleChatId(p.chatId)) return jsonRes(res, 400, { ok: false, error: 'invalid_chat_id' }); + jsonRes(res, 200, { + chatId: p.chatId, + listener: getMessageListenerConfig(cachedLarkAppId, p.chatId), + maxPromptBytes: MAX_MESSAGE_LISTENER_PROMPT_BYTES, + }); +}); + +ipcRoute('PUT', '/api/message-listeners/:chatId', async (req, res, p) => { + if (!cachedLarkAppId) return jsonRes(res, 503, { error: 'larkAppId_not_set' }); + if (!isValidRoleChatId(p.chatId)) return jsonRes(res, 400, { ok: false, error: 'invalid_chat_id' }); + let body: unknown; + try { body = await readJsonBody(req); } catch { return jsonRes(res, 400, { ok: false, error: 'bad_json' }); } + const update = sanitizeMessageListenerUpdate(body); + if (!update) return jsonRes(res, 400, { ok: false, error: 'invalid_listener' }); + const validation = validateMessageListenerUpdate(update); + if (!validation.ok) return jsonRes(res, 400, { ok: false, error: validation.reason }); + if (update.prompt && Buffer.byteLength(update.prompt, 'utf-8') > MAX_MESSAGE_LISTENER_PROMPT_BYTES) { + return jsonRes(res, 400, { ok: false, error: 'prompt_too_large' }); + } + const result = await updateMessageListenerConfig(cachedLarkAppId, p.chatId, update); + if (!result.ok) return jsonRes(res, ['prompt_required', 'sender_required'].includes(result.reason) ? 400 : 500, { ok: false, error: result.reason }); + jsonRes(res, 200, { ok: true, listener: result.listener }); +}); + +function dashboardHistoryMessageSender(message: any): { senderOpenId?: string; senderName?: string; senderTypeRaw?: string } { + const sender = message?.sender ?? {}; + const senderId = sender.id ?? sender.open_id ?? sender.user_id ?? sender.app_id + ?? message?.sender_id?.open_id ?? message?.sender_id?.user_id ?? message?.sender_id?.app_id; + const senderName = sender.sender_name ?? sender.name ?? sender.user_name ?? message?.sender_name; + const senderIdType = sender.id_type ?? sender.sender_id_type; + const senderTypeRaw = sender.sender_type ?? message?.sender_type ?? (senderIdType === 'app_id' ? 'app' : undefined); + return { + senderOpenId: typeof senderId === 'string' ? senderId : undefined, + senderName: typeof senderName === 'string' && senderName.trim() ? senderName.trim() : undefined, + senderTypeRaw: typeof senderTypeRaw === 'string' ? senderTypeRaw : undefined, + }; +} + +function dashboardMessageCreateTimeMs(message: any): number | undefined { + const value = Number(message?.create_time ?? message?.createTime); + return Number.isFinite(value) ? value : undefined; +} + +async function readMessageListenerPreviewRequest(req: IncomingMessage): Promise< + | { ok: true; listener: NonNullable>; limit: number } + | { ok: false; status: number; error: string } +> { + let body: unknown; + try { body = await readJsonBody(req); } catch { return { ok: false, status: 400, error: 'bad_json' }; } + const raw = body && typeof body === 'object' && !Array.isArray(body) ? body as Record : {}; + const listener = sanitizeMessageListenerUpdate(raw.listener ?? raw); + if (!listener) return { ok: false, status: 400, error: 'invalid_listener' }; + const validation = validateMessageListenerUpdate(listener); + if (!validation.ok) return { ok: false, status: 400, error: validation.reason }; + if (listener.prompt && Buffer.byteLength(listener.prompt, 'utf-8') > MAX_MESSAGE_LISTENER_PROMPT_BYTES) { + return { ok: false, status: 400, error: 'prompt_too_large' }; + } + return { ok: true, listener, limit: normalizeMessageListenerPreviewLimit(raw.limit) }; +} + +async function collectMessageListenerPreviewMatches( + larkAppId: string, + chatId: string, + listener: NonNullable>, + limit: number, +): Promise { + const bot = getBot(larkAppId); + const previewListener: MessageListenerConfig = { + enabled: true, + ...(listener.name ? { name: listener.name } : {}), + ...(listener.replyCardTitle ? { replyCardTitle: listener.replyCardTitle } : {}), + ...(listener.workingDir ? { workingDir: listener.workingDir } : {}), + prompt: listener.prompt, + ...(listener.senderPolicy && Object.keys(listener.senderPolicy).length > 0 ? { senderPolicy: listener.senderPolicy } : {}), + ...(listener.messagePolicy ? { messagePolicy: { ...listener.messagePolicy, scope: 'top_level' } } : { messagePolicy: { scope: 'top_level' } }), + replyPolicy: { mode: 'thread', sessionMode: 'per_message' }, + }; + const previewBot = { + ...bot, + config: { + ...bot.config, + messageListeners: { + ...(bot.config.messageListeners ?? {}), + [chatId]: previewListener, + }, + }, + }; + const cutoff = Date.now() - MESSAGE_LISTENER_PREVIEW_WINDOW_MS; + const messages = await listChatMessagesUntil(larkAppId, chatId, { + pageSize: 50, + stopAfter: (message, seenCount) => { + const createdAt = dashboardMessageCreateTimeMs(message); + return seenCount >= Math.max(100, limit * 5) || + (Number.isFinite(createdAt) && (createdAt as number) < cutoff); + }, + }); + return previewMessageListenerMatches({ + bot: previewBot, + chatId, + messages, + limit, + senderForMessage: dashboardHistoryMessageSender, + }); +} + +function publicMessageListenerMatch(match: MessageListenerPreviewMatch): Record { + return { + messageId: match.messageId, + createTime: match.createTime, + messageText: match.messageText, + messageTitle: match.messageTitle, + msgType: match.msgType, + senderOpenId: match.senderOpenId, + senderName: match.senderName, + senderType: match.senderType, + }; +} + +ipcRoute('POST', '/api/message-listeners/:chatId/preview', async (req, res, p) => { + if (!cachedLarkAppId) return jsonRes(res, 503, { ok: false, error: 'larkAppId_not_set' }); + if (!isValidRoleChatId(p.chatId)) return jsonRes(res, 400, { ok: false, error: 'invalid_chat_id' }); + const parsed = await readMessageListenerPreviewRequest(req); + if (!parsed.ok) return jsonRes(res, parsed.status, { ok: false, error: parsed.error }); + try { + const matches = await collectMessageListenerPreviewMatches(cachedLarkAppId, p.chatId, parsed.listener, parsed.limit); + jsonRes(res, 200, { + ok: true, + requestedLimit: parsed.limit, + matches: matches.map(publicMessageListenerMatch), + }); + } catch (err) { + jsonRes(res, 502, { ok: false, error: err instanceof Error ? err.message : String(err) }); + } +}); + +ipcRoute('POST', '/api/message-listeners/:chatId/run-preview', async (req, res, p) => { + if (!cachedLarkAppId) return jsonRes(res, 503, { ok: false, error: 'larkAppId_not_set' }); + if (!isValidRoleChatId(p.chatId)) return jsonRes(res, 400, { ok: false, error: 'invalid_chat_id' }); + const activeSessions = getActiveSessionsRegistry(); + if (!activeSessions) return jsonRes(res, 503, { ok: false, error: 'active session registry unavailable' }); + const parsed = await readMessageListenerPreviewRequest(req); + if (!parsed.ok) return jsonRes(res, parsed.status, { ok: false, error: parsed.error }); + try { + const matches = await collectMessageListenerPreviewMatches(cachedLarkAppId, p.chatId, parsed.listener, parsed.limit); + const run = createMessageListenerRunPreview(cachedLarkAppId, p.chatId, matches.map(match => match.messageId)); + const results = []; + for (const match of matches) { + const triggerId = createMessageListenerRunPreviewTurnId(); + try { + const result = await triggerSessionTurn({ + source: { + type: 'ui', + connectorId: 'message-listener-preview', + requestId: `listener-preview:${match.messageId}`, + receivedAt: new Date().toISOString(), + }, + target: { + kind: 'turn', + botId: cachedLarkAppId, + chatId: p.chatId, + rootMessageId: match.messageId, + }, + envelope: { + format: 'message_listener', + sourceName: match.name || 'Message Listener Preview', + trusted: false, + payload: publicMessageListenerMatch(match), + rawText: match.messageText, + }, + instruction: renderMessageListenerInstruction(match), + presentation: { topicMessage: null }, + }, { larkAppId: cachedLarkAppId, activeSessions }, { stableTurnId: triggerId }); + const tracked = result.ok + ? markMessageListenerRunPreviewTriggered(run.runId, match.messageId, { + action: result.action, + sessionId: result.target?.sessionId, + triggerId: result.triggerId ?? triggerId, + }) + : markMessageListenerRunPreviewFailed(run.runId, { + messageId: match.messageId, + sessionId: result.target?.sessionId, + error: result.error, + }); + results.push(tracked ?? { + runId: run.runId, + messageId: match.messageId, + ok: result.ok, + state: result.ok ? 'triggered' : 'failed', + action: result.action, + sessionId: result.target?.sessionId, + triggerId: result.triggerId ?? triggerId, + error: result.error, + }); + } catch (err) { + const error = err instanceof Error ? err.message : String(err); + const tracked = markMessageListenerRunPreviewFailed(run.runId, { + messageId: match.messageId, + error, + }); + results.push(tracked ?? { + runId: run.runId, + messageId: match.messageId, + ok: false, + state: 'failed', + error, + }); + } + } + jsonRes(res, 200, { + ok: results.every(result => result.ok), + runId: run.runId, + requestedLimit: parsed.limit, + matches: matches.map(publicMessageListenerMatch), + results, + }); + } catch (err) { + jsonRes(res, 502, { ok: false, error: err instanceof Error ? err.message : String(err) }); + } +}); + +ipcRoute('GET', '/api/message-listeners/:chatId/run-preview/:runId', async (_req, res, p) => { + if (!cachedLarkAppId) return jsonRes(res, 503, { ok: false, error: 'larkAppId_not_set' }); + if (!isValidRoleChatId(p.chatId)) return jsonRes(res, 400, { ok: false, error: 'invalid_chat_id' }); + const run = getMessageListenerRunPreview(p.runId); + if (!run || run.larkAppId !== cachedLarkAppId || run.chatId !== p.chatId) { + return jsonRes(res, 404, { ok: false, error: 'not_found' }); + } + jsonRes(res, 200, { + ok: true, + runId: run.runId, + createdAt: run.createdAt, + updatedAt: run.updatedAt, + results: run.results, + }); +}); + +ipcRoute('DELETE', '/api/message-listeners/:chatId', async (_req, res, p) => { + if (!cachedLarkAppId) return jsonRes(res, 503, { error: 'larkAppId_not_set' }); + if (!isValidRoleChatId(p.chatId)) return jsonRes(res, 400, { ok: false, error: 'invalid_chat_id' }); + const result = await updateMessageListenerConfig(cachedLarkAppId, p.chatId, { enabled: false, prompt: '' }); + if (!result.ok) return jsonRes(res, 500, { ok: false, error: result.reason }); + jsonRes(res, 200, { ok: true }); +}); + +ipcRoute('GET', '/api/groups/:chatId/members-display', async (_req, res, p) => { + if (!cachedLarkAppId) return jsonRes(res, 503, { error: 'larkAppId_not_set' }); + if (!isValidRoleChatId(p.chatId)) return jsonRes(res, 400, { ok: false, error: 'invalid_chat_id' }); + try { + const members = await listChatMemberDisplays(cachedLarkAppId, p.chatId); + jsonRes(res, 200, { members }); + } catch (err) { + jsonRes(res, 502, { ok: false, error: err instanceof Error ? err.message : String(err) }); + } +}); + // ─── Role profile management (dashboard) ────────────────────────────────── // Profiles are authoring/storage helpers only; applying one writes this bot's // entry into the selected chat role and does not alter runtime role layering. diff --git a/src/core/worker-pool.ts b/src/core/worker-pool.ts index b23f5b005..5093347d2 100644 --- a/src/core/worker-pool.ts +++ b/src/core/worker-pool.ts @@ -18,6 +18,11 @@ import { config } from '../config.js'; import { readGlobalConfig } from '../global-config.js'; import * as sessionStore from '../services/session-store.js'; import * as asyncTriggerStore from '../services/async-trigger-store.js'; +import { + markMessageListenerRunPreviewFailed, + markMessageListenerRunPreviewReplied, + markMessageListenerRunPreviewRunning, +} from '../services/message-listener-run-preview-store.js'; import { persistStreamCardState, rememberLastCliInput } from './session-manager.js'; import { fallbackTurnId, isSubstituteTurn } from './reply-target.js'; import { updateMessage, deleteMessage, sendEphemeralCard, sendUserMessage, addReaction, removeReaction, getMessageChatId, MessageWithdrawnError } from '../im/lark/client.js'; @@ -2563,6 +2568,9 @@ function setupWorkerHandlers( } if (recordDispatchInputCommit(ds.session, msg.turnId, workerGeneration)) { sessionStore.updateSession(ds.session); + if (msg.turnId.startsWith('mlrp_turn_')) { + markMessageListenerRunPreviewRunning(msg.turnId); + } } else { logger.warn(`[${t}] Ignored unbound input commit turn=${msg.turnId.slice(0, 16)}`); } @@ -3585,6 +3593,16 @@ function setupWorkerHandlers( break; } + case 'explicit_reply_observed': { + if (msg.turnId.startsWith('mlrp_turn_')) { + markMessageListenerRunPreviewReplied(msg.turnId, { + sessionId: ds.session.sessionId, + replyMessageId: msg.messageId, + }); + } + break; + } + case 'user_notify': { logger.warn(`[${t}] Worker user_notify: ${msg.message}`); emitSessionLifecycleHook(ds, 'session.requires_attention', { @@ -3627,6 +3645,12 @@ function setupWorkerHandlers( // never let a projection/store failure crash the worker IPC loop. logger.error(`[${t}] Failed to persist turn_terminal for ${msg.turnId.substring(0, 8)}: ${err.message}`); } + if (msg.turnId.startsWith('mlrp_turn_') && msg.status !== 'completed') { + markMessageListenerRunPreviewFailed(msg.turnId, { + sessionId: msg.sessionId, + error: msg.errorCode ?? msg.status, + }); + } try { await cb.onDeferredScheduleTurnSettled?.(ds, { turnId: msg.turnId, source: 'terminal' }); } catch (err: any) { @@ -3742,6 +3766,9 @@ function setupWorkerHandlers( clearUsageLimitState(ds); if (ds.lastScreenStatus === 'limited') ds.lastScreenStatus = 'idle'; } + if (msg.turnId.startsWith('mlrp_turn_')) { + markMessageListenerRunPreviewRunning(msg.turnId); + } deliverFinalOutput(ds, msg, t, 0); break; } @@ -4271,6 +4298,12 @@ function deliverFinalOutput( : undefined, ); recordPrimaryOutput(messageId); + if (msg.turnId.startsWith('mlrp_turn_')) { + markMessageListenerRunPreviewReplied(msg.turnId, { + sessionId: ds.session.sessionId, + replyMessageId: messageId, + }); + } if (preparedListenerReply?.kind === 'send' || preparedListenerReply?.kind === 'succeeded') { finishVcMeetingImReply(config.session.dataDir, preparedListenerReply.ref, messageId); } @@ -4288,6 +4321,12 @@ function deliverFinalOutput( const next = attempt + 1; if (next >= FINAL_OUTPUT_RETRY_BACKOFF_MS.length) { logger.error(`[${t}] Bridge final_output gave up after ${next} attempts (turn ${msg.turnId.substring(0, 8)}): ${err.message}`); + if (msg.turnId.startsWith('mlrp_turn_')) { + markMessageListenerRunPreviewFailed(msg.turnId, { + sessionId: ds.session.sessionId, + error: err.message, + }); + } // Don't commit the dedup marker — leave room for any future // retransmit (e.g. daemon restart that re-fires the IPC). return; diff --git a/src/daemon.ts b/src/daemon.ts index 4b3d9835f..078131ba3 100644 --- a/src/daemon.ts +++ b/src/daemon.ts @@ -175,6 +175,7 @@ import { extractBotmuxLarkNativeSessionTitlePrompt, } from './core/session-title.js'; import { settleDeferredScheduleRun } from './core/deferred-schedule-settlement.js'; +import { renderMessageListenerPrompt } from './services/message-listener.js'; import { sweepOrphanSandboxes } from './adapters/backend/sandbox.js'; import { TmuxBackend } from './adapters/backend/tmux-backend.js'; import { HerdrBackend } from './adapters/backend/herdr-backend.js'; @@ -2856,7 +2857,9 @@ export async function enforceMessageQuotaForCliInput( senderUnionId?: string, memberUnionId?: string, chatType?: 'group' | 'p2p', + opts?: { listenerAuthorized?: boolean }, ): Promise { + if (opts?.listenerAuthorized) return true; // senderUnionId(bot-locked)让 evaluateTalk 认出跨部署团队 peer bot(teamBot 腿); // memberUnionId(可为真人 union)走 teamMember 腿——否则外部闸门/群闸门放进来的 // 团队 bot 或团队成员消息会在这里复查处被静默丢弃(#332 端到端断点,人腿同理)。 @@ -14735,7 +14738,22 @@ async function resolvePinnedWorkingDir(ctx: { chatId: string; chatType: 'group' | 'p2p'; larkAppId: string; + listenerWorkingDir?: string; }) { + if (ctx.listenerWorkingDir) { + const resolved = expandHome(ctx.listenerWorkingDir); + try { + if (statSync(resolved).isDirectory()) { + return { + pinnedWorkingDir: resolved, + oncallEntry: undefined, + inheritedFrom: null, + pinnedFromBotDefault: false, + }; + } + } catch { /* fall through to normal resolution */ } + logger.warn(`[message-listener:${ctx.larkAppId}] listener workingDir invalid (${resolved}); falling back to normal resolution`); + } let oncallEntry = findOncallChat(ctx.larkAppId, ctx.chatId); if (!oncallEntry) { oncallEntry = await maybeAutoBindDefaultOncall(ctx.larkAppId, ctx.chatId, ctx.chatType); @@ -14991,7 +15009,7 @@ function mergeVcMeetingApplicationContext( } async function handleNewTopic(data: any, ctx: RoutingContext): Promise { - const { chatId, messageId, chatType, larkAppId, replyRootId, substituteTrigger } = ctx; + const { chatId, messageId, chatType, larkAppId, replyRootId, substituteTrigger, messageListener } = ctx; // scope/anchor are mutable here: `/t` / `/topic` may flip a 普通群 chat-scope // routing into thread-scope so the bot's first reply seeds a Lark thread. let scope = ctx.scope; @@ -15081,6 +15099,12 @@ async function handleNewTopic(data: any, ctx: RoutingContext): Promise { const teamTrustUnionId: string | undefined = (data.sender?.sender_type === 'app' || data.sender?.sender_type === 'bot') ? senderUnionId : undefined; const botCfg = getBot(larkAppId).config; + const listenerPrompt = messageListener ? renderMessageListenerPrompt(messageListener) : undefined; + if (listenerPrompt) { + content = listenerPrompt; + parsed.content = listenerPrompt; + cmdContent = listenerPrompt; + } logger.info(`New session: "${content.substring(0, 60)}" (scope=${scope}, anchor=${anchor.substring(0, 12)}, resources: ${resources.length}, active: ${getActiveCount()}, messageId: ${messageId}, chatId: ${chatId})`); emitHookEvent('topic.new', { larkAppId, @@ -15298,7 +15322,9 @@ async function handleNewTopic(data: any, ctx: RoutingContext): Promise { } } - if (!await enforceMessageQuotaForCliInput(larkAppId, chatId, senderOpenId, messageId, anchor, teamTrustUnionId, senderUnionId, chatType)) { + if (!await enforceMessageQuotaForCliInput(larkAppId, chatId, senderOpenId, messageId, anchor, teamTrustUnionId, senderUnionId, chatType, { + listenerAuthorized: !!messageListener, + })) { return; } @@ -15344,7 +15370,14 @@ async function handleNewTopic(data: any, ctx: RoutingContext): Promise { // Pin the working dir via the layered oncall / inherit / default lookup // (auto-binds a defaultOncall chat as a side effect). Shared with the // first-message `/repo` command branch so both paths stay consistent. - const { pinnedWorkingDir, oncallEntry, inheritedFrom, pinnedFromBotDefault } = await resolvePinnedWorkingDir({ scope, anchor, chatId, chatType, larkAppId }); + const { pinnedWorkingDir, oncallEntry, inheritedFrom, pinnedFromBotDefault } = await resolvePinnedWorkingDir({ + scope, + anchor, + chatId, + chatType, + larkAppId, + listenerWorkingDir: messageListener?.workingDir, + }); // Auto-worktree: register PENDING (router buffers concurrent msgs, no force-fork) // and build the worktree off the critical path (willAutoWorktree / runAutoWorktreeCommit). const autoWt = willAutoWorktree(larkAppId, pinnedWorkingDir, pinnedFromBotDefault); @@ -15359,7 +15392,8 @@ async function handleNewTopic(data: any, ctx: RoutingContext): Promise { // For chat-scope, rootMessageId stores the seed message_id (audit only); // routing keys off chatId via sessionAnchorId(), so any value works. const rootIdForStore = scope === 'thread' ? anchor : messageId; - const session = sessionStore.createSession(chatId, rootIdForStore, parsed.content.substring(0, 50), chatType); + const initialTurnTitle = (messageListener?.replyCardTitle ?? (ctx.forwardSeedData ? followupContent : content)).substring(0, 50); + const session = sessionStore.createSession(chatId, rootIdForStore, initialTurnTitle, chatType); const now = Date.now(); setDirectChatDisplayNameFromSender(session, chatType, newTopicSender); const groupChatName = await groupChatNamePromise; @@ -15419,7 +15453,7 @@ async function handleNewTopic(data: any, ctx: RoutingContext): Promise { pendingSubstituteControlCard: shouldSendSubstituteControlCard, pendingSender: newTopicSender, ownerOpenId: senderOpenId, - currentTurnTitle: (ctx.forwardSeedData ? followupContent : content).substring(0, 50), + currentTurnTitle: initialTurnTitle, workingDir: pinnedWorkingDir, }; if (pinnedWorkingDir) { diff --git a/src/dashboard.ts b/src/dashboard.ts index 28b849c3e..c594b47a5 100644 --- a/src/dashboard.ts +++ b/src/dashboard.ts @@ -2065,7 +2065,7 @@ async function buildGroupsMatrix(): Promise<{ chats: any[]; bots: any[] }> { if (!r.ok) return; const j = await r.json() as { chats?: any[] }; for (const c of j.chats ?? []) { - const { oncallChat, firstSeenAt, hasRole, observedBotNames, ...chatBase } = c; + const { oncallChat, firstSeenAt, hasRole, hasMessageListener, observedBotNames, ...chatBase } = c; const cur = out.get(c.chatId) ?? { ...chatBase, memberBots: [] as any[], @@ -2082,6 +2082,7 @@ async function buildGroupsMatrix(): Promise<{ chats: any[]; bots: any[] }> { inChat: true, oncallChat: oncallChat ?? null, hasRole: hasRole ?? false, + hasMessageListener: hasMessageListener ?? false, }); if (typeof firstSeenAt === 'number') { cur._firstSeenAt = cur._firstSeenAt === null @@ -2096,7 +2097,7 @@ async function buildGroupsMatrix(): Promise<{ chats: any[]; bots: any[] }> { const present = new Set(c.memberBots.map((mb: any) => mb.larkAppId)); for (const b of onlineBots) { if (!present.has(b.larkAppId)) { - c.memberBots.push({ larkAppId: b.larkAppId, botName: b.botName, cliId: b.cliId, inChat: false, oncallChat: null, hasRole: false }); + c.memberBots.push({ larkAppId: b.larkAppId, botName: b.botName, cliId: b.cliId, inChat: false, oncallChat: null, hasRole: false, hasMessageListener: false }); } } } @@ -3910,6 +3911,86 @@ const server = createServer(async (req, res) => { } } + // ─── Message listeners (proxy to daemon) ─────────────────────────────── + // GET /api/message-listeners/:larkAppId/:chatId + // PUT /api/message-listeners/:larkAppId/:chatId + // DELETE /api/message-listeners/:larkAppId/:chatId + // POST /api/message-listeners/:larkAppId/:chatId/(preview|run-preview) + // GET /api/message-listeners/:larkAppId/:chatId/run-preview/:runId + let mMessageListener: RegExpMatchArray | null; + if ((mMessageListener = url.pathname.match(/^\/api\/message-listeners\/([^/]+)\/([^/]+)\/run-preview\/([^/]+)$/))) { + const larkAppId = decodeURIComponent(mMessageListener[1]); + const chatId = decodeURIComponent(mMessageListener[2]); + const runId = decodeURIComponent(mMessageListener[3]); + if (req.method === 'GET') { + const upstream = await proxyToDaemon( + larkAppId, + `/api/message-listeners/${encodeURIComponent(chatId)}/run-preview/${encodeURIComponent(runId)}`, + { method: 'GET' }, + ); + res.writeHead(upstream.status, { 'content-type': 'application/json' }); + res.end(await upstream.text()); + return; + } + } + if ((mMessageListener = url.pathname.match(/^\/api\/message-listeners\/([^/]+)\/([^/]+)\/(preview|run-preview)$/))) { + const larkAppId = decodeURIComponent(mMessageListener[1]); + const chatId = decodeURIComponent(mMessageListener[2]); + const op = mMessageListener[3]; + if (req.method === 'POST') { + const chunks: Buffer[] = []; + for await (const c of req) chunks.push(c as Buffer); + const raw = Buffer.concat(chunks).toString('utf8') || '{}'; + const upstream = await proxyToDaemon(larkAppId, `/api/message-listeners/${encodeURIComponent(chatId)}/${op}`, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: raw, + }); + res.writeHead(upstream.status, { 'content-type': 'application/json' }); + res.end(await upstream.text()); + return; + } + } + if ((mMessageListener = url.pathname.match(/^\/api\/message-listeners\/([^/]+)\/([^/]+)$/))) { + const larkAppId = decodeURIComponent(mMessageListener[1]); + const chatId = decodeURIComponent(mMessageListener[2]); + if (req.method === 'GET') { + const upstream = await proxyToDaemon(larkAppId, `/api/message-listeners/${encodeURIComponent(chatId)}`, { method: 'GET' }); + res.writeHead(upstream.status, { 'content-type': 'application/json' }); + res.end(await upstream.text()); + return; + } + if (req.method === 'PUT') { + const chunks: Buffer[] = []; + for await (const c of req) chunks.push(c as Buffer); + const raw = Buffer.concat(chunks).toString('utf8') || '{}'; + const upstream = await proxyToDaemon(larkAppId, `/api/message-listeners/${encodeURIComponent(chatId)}`, { + method: 'PUT', + headers: { 'content-type': 'application/json' }, + body: raw, + }); + res.writeHead(upstream.status, { 'content-type': 'application/json' }); + res.end(await upstream.text()); + return; + } + if (req.method === 'DELETE') { + const upstream = await proxyToDaemon(larkAppId, `/api/message-listeners/${encodeURIComponent(chatId)}`, { method: 'DELETE' }); + res.writeHead(upstream.status, { 'content-type': 'application/json' }); + res.end(await upstream.text()); + return; + } + } + + let mGroupMembersDisplay: RegExpMatchArray | null; + if (req.method === 'GET' && (mGroupMembersDisplay = url.pathname.match(/^\/api\/groups\/([^/]+)\/([^/]+)\/members-display$/))) { + const larkAppId = decodeURIComponent(mGroupMembersDisplay[1]); + const chatId = decodeURIComponent(mGroupMembersDisplay[2]); + const upstream = await proxyToDaemon(larkAppId, `/api/groups/${encodeURIComponent(chatId)}/members-display`, { method: 'GET' }); + res.writeHead(upstream.status, { 'content-type': 'application/json' }); + res.end(await upstream.text()); + return; + } + // ─── Profiles (aggregate/proxy to daemon) ───────────────────────────── // ─── 会议角色预设(私有 API:不在 PUBLIC_READ_PATHS,未认证已被 401) ─── if (url.pathname === '/api/vc-meeting/consumer-profiles') { diff --git a/src/dashboard/web/i18n.ts b/src/dashboard/web/i18n.ts index b2d848b63..48ba7c35e 100644 --- a/src/dashboard/web/i18n.ts +++ b/src/dashboard/web/i18n.ts @@ -1958,6 +1958,62 @@ const zh: DashboardMessages = { 'roles.saveFailed': '保存失败,请重试', 'roles.tabGroups': '按群组', 'roles.tabProfiles': 'Profiles', + 'roles.editorTabs': '角色配置', + 'roles.roleTab': '角色提示词', + 'roles.listenerTab': '消息监听', + 'roles.listenerBadge': '监听', + 'roles.listenerEnabled': '启用群消息监听', + 'roles.listenerName': '监听名称', + 'roles.listenerNamePlaceholder': '例如:告警监听', + 'roles.listenerReplyCardTitle': '回复卡片标题', + 'roles.listenerReplyCardTitlePlaceholder': '留空使用默认标题', + 'roles.listenerWorkingDir': '工作目录', + 'roles.listenerWorkingDirPlaceholder': '留空使用 Bot 默认工作目录', + 'roles.listenerSenderTypes': '发送者类型', + 'roles.listenerSenderUser': '用户', + 'roles.listenerSenderBot': '机器人', + 'roles.listenerMessageTypes': '消息类型', + 'roles.listenerExcludeSelf': '排除当前机器人', + 'roles.listenerPrompt': '监听提示词', + 'roles.listenerPromptPlaceholder': '描述哪些群消息需要处理,以及如何回复。命中的消息会在原消息下方新建话题回复。', + 'roles.listenerPromptRequired': '启用监听时提示词不能为空', + 'roles.listenerSenderRequired': '启用监听时至少选择一个发送者', + 'roles.listenerFilters': '发送者过滤', + 'roles.listenerMembers': '成员过滤', + 'roles.listenerBots': '机器人过滤', + 'roles.listenerMembersHint': '显示名仅用于选择,配置保存 open_id', + 'roles.listenerMembersEmpty': '暂无成员数据', + 'roles.listenerBotsEmpty': '暂无机器人数据', + 'roles.listenerSearchPlaceholder': '按名称或 open_id 搜索', + 'roles.listenerClearSearch': '清空搜索', + 'roles.listenerSelectVisible': '全选 {count} 项', + 'roles.listenerClearVisible': '取消全选 {count} 项', + 'roles.listenerSelectedCount': '已选择 {count} 项', + 'roles.listenerFilteredCount': '当前筛选 {count} 项', + 'roles.listenerMixedMode': '混合状态', + 'roles.listenerBulkActions': '批量设置监听状态', + 'roles.listenerTargetListen': '监听', + 'roles.listenerTargetIgnore': '不监听', + 'roles.listenerConfirmDelete': '确认删除该 Bot 在此群的消息监听配置?', + 'roles.listenerScopeHelp': '监听只处理群聊顶层消息;显式 @ 当前 Bot 的消息仍走普通 @ 路由;默认不处理已有线程里的普通回复。', + 'roles.listenerPreviewTitle': '预览与试运行', + 'roles.listenerPreviewHint': '预览最近 24 小时内符合当前条件的群消息;试运行会按真实链路创建会话并回复。', + 'roles.listenerPreviewCount': '数量', + 'roles.listenerPreviewButton': '预览', + 'roles.listenerRunPreviewButton': '试运行', + 'roles.listenerPreviewEmpty': '尚未预览。', + 'roles.listenerPreviewNoMatches': '最近 24 小时内没有符合条件的内容。', + 'roles.listenerPreviewSummary': '命中 {count} 条消息', + 'roles.listenerRunPreviewSummary': '试运行 {count} 条:已触发 {triggered} 条,执行中 {running} 条,已回复 {replied} 条,回复失败 {failed} 条', + 'roles.listenerPreviewFailed': '预览失败,请重试', + 'roles.listenerPreviewSender': '发送者', + 'roles.listenerPreviewTime': '发送时间', + 'roles.listenerPreviewMessageTitle': '消息标题', + 'roles.listenerPreviewContent': '消息内容', + 'roles.listenerRunPreviewState.triggered': '已触发', + 'roles.listenerRunPreviewState.running': '执行中', + 'roles.listenerRunPreviewState.replied': '已回复', + 'roles.listenerRunPreviewState.failed': '回复失败', 'roles.profileIdPlaceholder': 'profile id,如 collab-main', 'roles.profileIdInvalid': 'Profile id 只能包含字母、数字、点、下划线、短横线,最长 64,且不能是 . 或 ..', 'roles.openProfile': '打开', @@ -3956,6 +4012,62 @@ const en: DashboardMessages = { 'roles.saveFailed': 'Save failed, please retry', 'roles.tabGroups': 'By Group', 'roles.tabProfiles': 'Profiles', + 'roles.editorTabs': 'Role settings', + 'roles.roleTab': 'Role Prompt', + 'roles.listenerTab': 'Message Listener', + 'roles.listenerBadge': 'Listener', + 'roles.listenerEnabled': 'Enable group message listener', + 'roles.listenerName': 'Listener name', + 'roles.listenerNamePlaceholder': 'e.g. Alert listener', + 'roles.listenerReplyCardTitle': 'Reply card title', + 'roles.listenerReplyCardTitlePlaceholder': 'Leave empty to use the default title', + 'roles.listenerWorkingDir': 'Working directory', + 'roles.listenerWorkingDirPlaceholder': 'Leave empty to use the bot default', + 'roles.listenerSenderTypes': 'Sender types', + 'roles.listenerSenderUser': 'Users', + 'roles.listenerSenderBot': 'Bots', + 'roles.listenerMessageTypes': 'Message types', + 'roles.listenerExcludeSelf': 'Exclude this bot', + 'roles.listenerPrompt': 'Listener prompt', + 'roles.listenerPromptPlaceholder': 'Describe which group messages should be handled and how to reply. Matches reply as a topic under the original message.', + 'roles.listenerPromptRequired': 'Prompt is required when listener is enabled', + 'roles.listenerSenderRequired': 'Select at least one sender before enabling the listener', + 'roles.listenerFilters': 'Sender filters', + 'roles.listenerMembers': 'Member filters', + 'roles.listenerBots': 'Bot filters', + 'roles.listenerMembersHint': 'Display names are only for picking; open_id is saved', + 'roles.listenerMembersEmpty': 'No members loaded', + 'roles.listenerBotsEmpty': 'No bots loaded', + 'roles.listenerSearchPlaceholder': 'Search by name or open_id', + 'roles.listenerClearSearch': 'Clear search', + 'roles.listenerSelectVisible': 'Select {count}', + 'roles.listenerClearVisible': 'Clear {count}', + 'roles.listenerSelectedCount': '{count} selected', + 'roles.listenerFilteredCount': '{count} filtered', + 'roles.listenerMixedMode': 'mixed', + 'roles.listenerBulkActions': 'Bulk listening state', + 'roles.listenerTargetListen': 'Listen', + 'roles.listenerTargetIgnore': 'Ignore', + 'roles.listenerConfirmDelete': 'Delete this bot\'s message listener config for this group?', + 'roles.listenerScopeHelp': 'The listener only handles top-level group messages. Messages that explicitly mention this bot still use normal mention routing. Existing thread replies are not handled by default.', + 'roles.listenerPreviewTitle': 'Preview and Test Run', + 'roles.listenerPreviewHint': 'Preview group messages from the last 24 hours that match the current filters. Test run creates real sessions and replies through the normal path.', + 'roles.listenerPreviewCount': 'Count', + 'roles.listenerPreviewButton': 'Preview', + 'roles.listenerRunPreviewButton': 'Test run', + 'roles.listenerPreviewEmpty': 'No preview yet.', + 'roles.listenerPreviewNoMatches': 'No messages from the last 24 hours match these filters.', + 'roles.listenerPreviewSummary': '{count} matching messages', + 'roles.listenerRunPreviewSummary': 'Test ran {count}: triggered {triggered}, running {running}, replied {replied}, failed {failed}', + 'roles.listenerPreviewFailed': 'Preview failed. Try again.', + 'roles.listenerPreviewSender': 'Sender', + 'roles.listenerPreviewTime': 'Sent at', + 'roles.listenerPreviewMessageTitle': 'Message title', + 'roles.listenerPreviewContent': 'Message content', + 'roles.listenerRunPreviewState.triggered': 'Triggered', + 'roles.listenerRunPreviewState.running': 'Running', + 'roles.listenerRunPreviewState.replied': 'Replied', + 'roles.listenerRunPreviewState.failed': 'Reply failed', 'roles.profileIdPlaceholder': 'profile id, e.g. collab-main', 'roles.profileIdInvalid': 'Profile id may contain letters, numbers, dots, underscores, and hyphens only; max 64; cannot be . or ..', 'roles.openProfile': 'Open', diff --git a/src/dashboard/web/listener-filters.ts b/src/dashboard/web/listener-filters.ts new file mode 100644 index 000000000..3b2a19e57 --- /dev/null +++ b/src/dashboard/web/listener-filters.ts @@ -0,0 +1,55 @@ +export type ListenerTargetState = 'listen' | 'ignore'; +export type ListenerTargetBulkState = ListenerTargetState | 'mixed'; + +export interface ListenerFilterTarget { + openId: string; + name: string; + memberType: 'user' | 'bot' | 'unknown'; +} + +export function filterListenerTargets(targets: T[], query: string): T[] { + const q = query.trim().toLowerCase(); + if (!q) return targets; + return targets.filter(target => + target.openId.toLowerCase().includes(q) + || target.name.toLowerCase().includes(q) + || target.memberType.toLowerCase().includes(q), + ); +} + +function unique(values: Iterable): string[] { + return [...new Set([...values].filter(Boolean))]; +} + +export function applyListenerFilterState(input: { + include: readonly string[]; + exclude: readonly string[]; + targetIds: readonly string[]; + listening: boolean; +}): { include: string[]; exclude: string[] } { + const targetIds = new Set(input.targetIds.filter(Boolean)); + const include = new Set(input.include); + for (const id of targetIds) { + include.delete(id); + if (input.listening) include.add(id); + } + return { + include: unique(include), + exclude: [], + }; +} + +export function listenerTargetStateFor(input: { + include: readonly string[]; + exclude: readonly string[]; + targetIds: readonly string[]; +}): ListenerTargetBulkState { + const targetIds = input.targetIds.filter(Boolean); + if (targetIds.length === 0) return 'ignore'; + const include = new Set(input.include); + const states = new Set(); + for (const id of targetIds) { + states.add(include.has(id) ? 'listen' : 'ignore'); + } + return states.size === 1 ? [...states][0] : 'mixed'; +} diff --git a/src/dashboard/web/roles-page.tsx b/src/dashboard/web/roles-page.tsx index 910a23c54..b77e12602 100644 --- a/src/dashboard/web/roles-page.tsx +++ b/src/dashboard/web/roles-page.tsx @@ -8,6 +8,11 @@ import { type SetStateAction, } from 'react'; import { DropdownMenu, Html, LoadingState, RefreshIconButton } from './dashboard-components.js'; +import { + applyListenerFilterState, + filterListenerTargets, + listenerTargetStateFor, +} from './listener-filters.js'; import { useT } from './react-hooks.js'; import { mountReactPage, type PageDisposer } from './react-mount.js'; import { @@ -21,26 +26,44 @@ import { botRoleCount, byteLength, deleteProfileEntry, + deleteMessageListener, deleteRole, entryForBot, filterRoleGroups, filterRoleProfiles, + formatListenerPreviewTime, hashChatId, isValidProfileId, + loadGroupMemberDisplays, loadGroups, + loadMessageListenerRunPreviewStatus, + loadMessageListener, loadProfileEntries, loadProfileEntry, loadProfiles, loadRole, loadRoleProfileContext, + MAX_MESSAGE_LISTENER_PROMPT_BYTES, MAX_ROLE_BYTES, + MESSAGE_LISTENER_WARN_BYTES, + DEFAULT_MESSAGE_LISTENER_PREVIEW_LIMIT, + MAX_MESSAGE_LISTENER_PREVIEW_LIMIT, roleKey, ROLE_WARN_BYTES, saveInjectMode, + saveMessageListener, saveProfileEntry, saveRole, + previewMessageListener, + runMessageListenerPreview, type DashboardBot, type GroupInfo, + type GroupMemberDisplay, + type MessageListenerData, + type MessageListenerPreviewItem, + type MessageListenerPreviewResponse, + type MessageListenerRunPreviewResult, + type MessageListenerRunPreviewState, type RoleData, type RoleInjectMode, type RoleProfileApplyResult, @@ -51,13 +74,109 @@ import { import { botAvatarHtml, loadNameMaps } from './ui.js'; type RolesTab = 'groups' | 'profiles'; +type GroupEditorSection = 'role' | 'listener'; +type ListenerTargetTab = 'members' | 'bots'; +type SenderTypeOption = 'user' | 'bot'; type Translator = ReturnType; +const LISTENER_MESSAGE_TYPES = ['text', 'post', 'image', 'interactive'] as const; + type FlashState = { text: string; isError?: boolean; id: number } | null; type ApplyStatus = | { kind: 'idle' } | { kind: 'text'; text: string } | { kind: 'results'; preview: boolean; results: RoleProfileApplyResult[] }; +type ListenerPreviewStatus = + | { kind: 'idle' } + | { kind: 'loading'; mode: 'preview' | 'run' } + | { kind: 'result'; response: MessageListenerPreviewResponse; mode: 'preview' | 'run' } + | { kind: 'error'; text: string }; + +const DEFAULT_LISTENER: MessageListenerData = { + enabled: false, + prompt: '', + senderPolicy: { + mode: 'include_only', + includeSenderTypes: ['user'], + excludeSelf: true, + }, + messagePolicy: { + includeMsgTypes: [...LISTENER_MESSAGE_TYPES], + scope: 'top_level', + }, +}; + +function cloneListener(listener: MessageListenerData | null | undefined): MessageListenerData { + return { + enabled: listener?.enabled === true, + name: listener?.name ?? '', + replyCardTitle: listener?.replyCardTitle ?? '', + workingDir: listener?.workingDir ?? '', + prompt: listener?.prompt ?? '', + senderPolicy: { + mode: 'include_only', + includeSenderOpenIds: [...(listener?.senderPolicy?.includeSenderOpenIds ?? [])], + excludeSenderOpenIds: [], + includeSenderTypes: [...(listener?.senderPolicy?.includeSenderTypes ?? DEFAULT_LISTENER.senderPolicy?.includeSenderTypes ?? [])], + excludeSenderTypes: [...(listener?.senderPolicy?.excludeSenderTypes ?? [])], + excludeSelf: listener?.senderPolicy?.excludeSelf !== false, + }, + messagePolicy: { + includeMsgTypes: [...(listener?.messagePolicy?.includeMsgTypes ?? DEFAULT_LISTENER.messagePolicy?.includeMsgTypes ?? [])], + scope: 'top_level', + }, + }; +} + +function listenerHasConfig(listener: MessageListenerData | null): boolean { + return listener?.enabled === true && listener.prompt.trim().length > 0; +} + +function groupHasAnyRoleOrListener(group: GroupInfo): boolean { + return group.memberBots.some(bot => bot.inChat && (bot.hasRole || bot.hasMessageListener)); +} + +function memberDisplayName(member: GroupMemberDisplay | undefined, openId: string): string { + return member?.name || openId; +} + +function listenerSenderTypeMatches(member: GroupMemberDisplay, listener: MessageListenerData): boolean { + const includeTypes = new Set(listener.senderPolicy?.includeSenderTypes ?? DEFAULT_LISTENER.senderPolicy?.includeSenderTypes ?? []); + if (includeTypes.size === 0) return true; + return (member.memberType === 'user' || member.memberType === 'bot') && includeTypes.has(member.memberType); +} + +function mergeListenerRunPreviewResults( + current: MessageListenerRunPreviewResult[] | undefined, + next: MessageListenerRunPreviewResult[], +): MessageListenerRunPreviewResult[] { + const merged = new Map(); + for (const result of current ?? []) merged.set(result.messageId, result); + for (const result of next) merged.set(result.messageId, { ...merged.get(result.messageId), ...result }); + return [...merged.values()]; +} + +function listenerRunPreviewStateClass(state: MessageListenerRunPreviewState | undefined, ok: boolean): string { + if (!ok || state === 'failed') return 'error'; + if (state === 'replied') return 'ok'; + if (state === 'running') return 'running'; + return 'triggered'; +} + +function listenerForEditor(listener: MessageListenerData | null | undefined, members: GroupMemberDisplay[] = []): MessageListenerData { + const next = cloneListener(listener ?? DEFAULT_LISTENER); + if (listener?.senderPolicy?.mode !== 'all_except_excluded') return next; + const excluded = new Set(listener.senderPolicy.excludeSenderOpenIds ?? []); + next.senderPolicy = { + ...(next.senderPolicy ?? {}), + mode: 'include_only', + includeSenderOpenIds: members + .filter(member => listenerSenderTypeMatches(member, next) && !excluded.has(member.openId)) + .map(member => member.openId), + excludeSenderOpenIds: [], + }; + return next; +} function useAliveRef() { const alive = useRef(true); @@ -123,6 +242,19 @@ function RolesPage(props: { tab: RolesTab }) { const [injectSaving, setInjectSaving] = useState(false); const [roleFlash, setRoleFlash] = useState(null); const [injectFlash, setInjectFlash] = useState(null); + const [groupEditorSection, setGroupEditorSection] = useState('role'); + const [selectedListener, setSelectedListener] = useState(null); + const [editingListener, setEditingListener] = useState(() => cloneListener(DEFAULT_LISTENER)); + const [listenerMembers, setListenerMembers] = useState([]); + const [listenerLoading, setListenerLoading] = useState(false); + const [listenerMembersLoading, setListenerMembersLoading] = useState(false); + const [listenerSaving, setListenerSaving] = useState(false); + const [listenerDeleting, setListenerDeleting] = useState(false); + const [listenerFlash, setListenerFlash] = useState(null); + const [listenerPreviewLimit, setListenerPreviewLimit] = useState(DEFAULT_MESSAGE_LISTENER_PREVIEW_LIMIT); + const [listenerPreviewStatus, setListenerPreviewStatus] = useState({ kind: 'idle' }); + const listenerRunPollRef = useRef<{ runId: string; token: number } | null>(null); + const listenerRunPollToken = useRef(0); const [selectedProfileId, setSelectedProfileId] = useState(null); const [selectedProfileBotId, setSelectedProfileBotId] = useState(null); const [profileEntries, setProfileEntries] = useState([]); @@ -163,6 +295,12 @@ function RolesPage(props: { tab: RolesTab }) { ); const roleByteLen = byteLength(editingContent); const profileByteLen = byteLength(profileEditingContent); + const listenerPromptByteLen = byteLength(editingListener.prompt); + const listenerMemberById = useMemo(() => { + const map = new Map(); + for (const member of listenerMembers) map.set(member.openId, member); + return map; + }, [listenerMembers]); const flash = useCallback((setter: Dispatch>, text: string, isError = false) => { const id = Date.now() + Math.random(); @@ -200,6 +338,33 @@ function RolesPage(props: { tab: RolesTab }) { return nextProfiles; }, [alive]); + const applyLoadedListener = useCallback((listener: MessageListenerData | null, members: GroupMemberDisplay[] = []) => { + const next = listenerForEditor(listener ?? DEFAULT_LISTENER, members); + setSelectedListener(listener); + setEditingListener(next); + }, []); + + const loadListenerForSelection = useCallback(async (botId: string, groupId: string, serial: number) => { + setListenerLoading(true); + setListenerMembersLoading(true); + setListenerMembers([]); + try { + const [listener, members] = await Promise.all([ + loadMessageListener(botId, groupId).catch(() => cloneListener(DEFAULT_LISTENER)), + loadGroupMemberDisplays(botId, groupId).catch(() => [] as GroupMemberDisplay[]), + ]); + if (!alive.current || serial !== selectSerial.current) return; + applyLoadedListener(listenerHasConfig(listener) ? listener : null, members); + setListenerMembers(members); + setListenerFlash(null); + } finally { + if (alive.current && serial === selectSerial.current) { + setListenerLoading(false); + setListenerMembersLoading(false); + } + } + }, [alive, applyLoadedListener]); + const loadInitial = useCallback(async () => { setLoadingTree(true); setProfileListLoading(true); @@ -212,7 +377,7 @@ function RolesPage(props: { tab: RolesTab }) { await loadNameMaps(); if (!alive.current) return; - setExpandedGroups(new Set(snapshot.groups.filter(group => botRoleCount(group) > 0).map(group => group.chatId))); + setExpandedGroups(new Set(snapshot.groups.filter(groupHasAnyRoleOrListener).map(group => group.chatId))); if (props.tab === 'profiles') { const requestedChatId = hashChatId(); setSelectedApplyGroupId(current => { @@ -238,6 +403,11 @@ function RolesPage(props: { tab: RolesTab }) { void loadInitial(); }, [loadInitial]); + useEffect(() => () => { + listenerRunPollToken.current += 1; + listenerRunPollRef.current = null; + }, []); + useEffect(() => { if (selectedApplyGroupId || groups.length === 0) return; setSelectedApplyGroupId(groups[0].chatId); @@ -257,13 +427,16 @@ function RolesPage(props: { tab: RolesTab }) { const serial = ++selectSerial.current; setSelectedGroupId(groupId); setSelectedBotId(botId); + setRoleFlash(null); + setInjectFlash(null); + setListenerFlash(null); + applyLoadedListener(null); const role = await loadRole(botId, groupId); if (!alive.current || serial !== selectSerial.current) return; setSelectedRole(role); setEditingContent(role.content ?? ''); setEditingInjectMode(role.injectMode === 'once' ? 'once' : 'every'); - setRoleFlash(null); - setInjectFlash(null); + await loadListenerForSelection(botId, groupId, serial); } async function handleGroupRefresh(): Promise { @@ -271,11 +444,13 @@ function RolesPage(props: { tab: RolesTab }) { if (!alive.current) return; void refreshRoleContext(snapshot.groups, profiles); if (selectedGroupId && selectedBotId) { + const serial = ++selectSerial.current; const role = await loadRole(selectedBotId, selectedGroupId); - if (!alive.current) return; + if (!alive.current || serial !== selectSerial.current) return; setSelectedRole(role); setEditingContent(role.content ?? ''); setEditingInjectMode(role.injectMode === 'once' ? 'once' : 'every'); + await loadListenerForSelection(selectedBotId, selectedGroupId, serial); } } @@ -336,6 +511,258 @@ function RolesPage(props: { tab: RolesTab }) { } } + function updateEditingListener(patch: Partial): void { + setEditingListener(prev => ({ ...prev, ...patch })); + } + + function updateListenerSenderPolicy(patch: NonNullable): void { + setEditingListener(prev => ({ + ...prev, + senderPolicy: { + ...(prev.senderPolicy ?? {}), + ...patch, + }, + })); + } + + function updateListenerMessagePolicy(patch: NonNullable): void { + setEditingListener(prev => ({ + ...prev, + messagePolicy: { + ...(prev.messagePolicy ?? { scope: 'top_level' }), + ...patch, + scope: 'top_level', + }, + })); + } + + function toggleListenerSenderType(type: SenderTypeOption, checked: boolean): void { + setEditingListener(prev => { + const current = new Set(prev.senderPolicy?.includeSenderTypes ?? []); + if (checked) current.add(type); + else { + if (current.size <= 1 && current.has(type)) return prev; + current.delete(type); + } + return { + ...prev, + senderPolicy: { + ...(prev.senderPolicy ?? {}), + includeSenderTypes: [...current], + }, + }; + }); + } + + function toggleListenerMsgType(msgType: string, checked: boolean): void { + setEditingListener(prev => { + const current = new Set(prev.messagePolicy?.includeMsgTypes ?? []); + if (checked) current.add(msgType); + else { + if (current.size <= 1 && current.has(msgType)) return prev; + current.delete(msgType); + } + return { + ...prev, + messagePolicy: { + ...(prev.messagePolicy ?? { scope: 'top_level' }), + includeMsgTypes: [...current], + scope: 'top_level', + }, + }; + }); + } + + function setListenerTargetPolicy(openId: string, listening: boolean): void { + setEditingListener(prev => { + const next = applyListenerFilterState({ + include: prev.senderPolicy?.includeSenderOpenIds ?? [], + exclude: prev.senderPolicy?.excludeSenderOpenIds ?? [], + targetIds: [openId], + listening, + }); + return { + ...prev, + senderPolicy: { + ...(prev.senderPolicy ?? {}), + mode: 'include_only', + includeSenderOpenIds: next.include, + excludeSenderOpenIds: next.exclude, + }, + }; + }); + } + + function setListenerTargetsPolicy(openIds: string[], listening: boolean): void { + setEditingListener(prev => { + const current = applyListenerFilterState({ + include: prev.senderPolicy?.includeSenderOpenIds ?? [], + exclude: prev.senderPolicy?.excludeSenderOpenIds ?? [], + targetIds: openIds, + listening, + }); + return { + ...prev, + senderPolicy: { + ...(prev.senderPolicy ?? {}), + mode: 'include_only', + includeSenderOpenIds: current.include, + excludeSenderOpenIds: current.exclude, + }, + }; + }); + } + + function listenerSavePayload(): MessageListenerData { + const senderPolicy = editingListener.senderPolicy ?? {}; + const messagePolicy = editingListener.messagePolicy ?? {}; + const includeSenderOpenIds = [...new Set(senderPolicy.includeSenderOpenIds ?? [])].filter(Boolean); + const includeSenderTypes = [...new Set(senderPolicy.includeSenderTypes ?? [])].filter((type): type is SenderTypeOption => type === 'user' || type === 'bot'); + const includeMsgTypes = [...new Set(messagePolicy.includeMsgTypes ?? [])].filter(Boolean); + return { + enabled: editingListener.enabled, + ...(editingListener.name?.trim() ? { name: editingListener.name.trim() } : {}), + ...(editingListener.replyCardTitle?.trim() ? { replyCardTitle: editingListener.replyCardTitle.trim() } : {}), + ...(editingListener.workingDir?.trim() ? { workingDir: editingListener.workingDir.trim() } : {}), + prompt: editingListener.prompt.trim(), + senderPolicy: { + ...(includeSenderOpenIds.length > 0 ? { includeSenderOpenIds } : {}), + mode: 'include_only', + ...(includeSenderTypes.length > 0 ? { includeSenderTypes } : {}), + excludeSelf: senderPolicy.excludeSelf !== false, + }, + messagePolicy: { + ...(includeMsgTypes.length > 0 ? { includeMsgTypes } : {}), + scope: 'top_level', + }, + }; + } + + function validateListenerForPreview(): MessageListenerData | null { + if (!editingListener.prompt.trim()) { + flash(setListenerFlash, tr('roles.listenerPromptRequired'), true); + return null; + } + if ((editingListener.senderPolicy?.includeSenderOpenIds?.length ?? 0) === 0) { + flash(setListenerFlash, tr('roles.listenerSenderRequired'), true); + return null; + } + return { + ...listenerSavePayload(), + enabled: true, + }; + } + + async function handleListenerPreview(run: boolean): Promise { + if (!selectedGroupId || !selectedBotId) return; + const payload = validateListenerForPreview(); + if (!payload) return; + const mode = run ? 'run' : 'preview'; + setListenerPreviewStatus({ kind: 'loading', mode }); + try { + const response = run + ? await runMessageListenerPreview(selectedBotId, selectedGroupId, payload, listenerPreviewLimit) + : await previewMessageListener(selectedBotId, selectedGroupId, payload, listenerPreviewLimit); + if (!alive.current) return; + setListenerPreviewStatus(response.ok + ? { kind: 'result', response, mode } + : { kind: 'error', text: response.error || tr('roles.listenerPreviewFailed') }); + if (response.ok && run && response.runId) { + startListenerRunPreviewPolling(response.runId); + } + } catch (err) { + if (!alive.current) return; + setListenerPreviewStatus({ kind: 'error', text: err instanceof Error ? err.message : tr('roles.listenerPreviewFailed') }); + } + } + + function startListenerRunPreviewPolling(runId: string): void { + if (!selectedGroupId || !selectedBotId) return; + const token = ++listenerRunPollToken.current; + listenerRunPollRef.current = { runId, token }; + const poll = async () => { + if (!alive.current) return; + const current = listenerRunPollRef.current; + if (!current || current.runId !== runId || current.token !== token || !selectedGroupId || !selectedBotId) return; + try { + const status = await loadMessageListenerRunPreviewStatus(selectedBotId, selectedGroupId, runId); + if (!alive.current) return; + if (status.ok && status.results) { + const nextResults = status.results; + setListenerPreviewStatus(previous => { + if (previous.kind !== 'result' || previous.mode !== 'run' || previous.response.runId !== runId) return previous; + return { + ...previous, + response: { + ...previous.response, + results: mergeListenerRunPreviewResults(previous.response.results, nextResults), + }, + }; + }); + if (nextResults.some(result => result.state === 'triggered' || result.state === 'running')) { + scheduleTimer(poll, 1500); + } else if (listenerRunPollRef.current?.runId === runId) { + listenerRunPollRef.current = null; + } + } + } catch { + scheduleTimer(poll, 2500); + } + }; + scheduleTimer(poll, 1500); + } + + async function handleSaveListener(): Promise { + if (!selectedGroupId || !selectedBotId) return; + if (!editingListener.enabled) { + await handleDeleteListener(false); + return; + } + if (!editingListener.prompt.trim()) { + flash(setListenerFlash, tr('roles.listenerPromptRequired'), true); + return; + } + if ((editingListener.senderPolicy?.includeSenderOpenIds?.length ?? 0) === 0) { + flash(setListenerFlash, tr('roles.listenerSenderRequired'), true); + return; + } + setListenerSaving(true); + try { + const ok = await saveMessageListener(selectedBotId, selectedGroupId, listenerSavePayload()); + if (!alive.current) return; + if (ok) { + const snapshot = await refreshGroups(); + if (!alive.current) return; + const listener = await loadMessageListener(selectedBotId, selectedGroupId); + if (!alive.current) return; + applyLoadedListener(listenerHasConfig(listener) ? listener : null); + void refreshRoleContext(snapshot.groups, profiles); + } + flash(setListenerFlash, ok ? tr('roles.saved') : tr('roles.saveFailed'), !ok); + } finally { + if (alive.current) setListenerSaving(false); + } + } + + async function handleDeleteListener(confirmFirst = true): Promise { + if (!selectedGroupId || !selectedBotId) return; + if (confirmFirst && !confirm(tr('roles.listenerConfirmDelete'))) return; + setListenerDeleting(true); + try { + const ok = await deleteMessageListener(selectedBotId, selectedGroupId); + if (!alive.current) return; + if (ok) { + const snapshot = await refreshGroups(); + if (!alive.current) return; + applyLoadedListener(null); + void refreshRoleContext(snapshot.groups, profiles); + } + flash(setListenerFlash, ok ? tr('roles.saved') : tr('roles.saveFailed'), !ok); + } finally { + if (alive.current) setListenerDeleting(false); + } + } + async function handleSelectProfile(profileId: string): Promise { const clean = profileId.trim(); if (!isValidProfileId(clean)) return; @@ -467,6 +894,10 @@ function RolesPage(props: { tab: RolesTab }) { } const roleSaveDisabled = roleSaving || roleByteLen > MAX_ROLE_BYTES || editingContent.trim().length === 0; + const listenerSaveDisabled = listenerSaving + || listenerDeleting + || listenerPromptByteLen > MAX_MESSAGE_LISTENER_PROMPT_BYTES + || (editingListener.enabled && editingListener.prompt.trim().length === 0); const profileSaveDisabled = profileSaving || profileByteLen > MAX_ROLE_BYTES || profileEditingContent.trim().length === 0; const isProfiles = props.tab === 'profiles'; const tabs = ( @@ -526,57 +957,128 @@ function RolesPage(props: { tab: RolesTab }) {
- - + {groupEditorSection === 'role' ? ( + <> + + + + ) : ( + <> + + + + )}
-
- {tr('roles.injectModeLabel')} - void handleInjectModeChange(mode === 'once' ? 'once' : 'every')} - /> - {tr('roles.injectModeHint')} - -
-