From 0c4547e1b67024207c2da902bc9aa3745313a1a5 Mon Sep 17 00:00:00 2001 From: "zhujinhao.zjh" Date: Tue, 25 Aug 2026 16:39:38 +0800 Subject: [PATCH 1/2] fix(listener): prevent daemon stalls during group polling Keep restored persistent sessions lazy by default, make dashboard bulk rows lightweight, and normalize/cancel REST chat-history polling so group listeners stay responsive on hosts with large session history. Co-authored-by: TRAE CLI --- src/config.ts | 11 +-- src/core/dashboard-ipc-server.ts | 14 +-- src/core/dashboard-rows.ts | 43 ++++++--- src/core/session-activity.ts | 4 +- src/core/session-manager.ts | 88 ++++++++++++------ src/im/lark/client.ts | 6 +- src/im/lark/event-dispatcher.ts | 93 ++++++++++++++++--- src/im/lark/message-parser.ts | 18 +++- test/dashboard-attention-signals.test.ts | 8 ++ test/dashboard-row-lightweight-source.test.ts | 33 +++++++ test/event-dispatcher.test.ts | 55 ++++++++++- test/message-parser.test.ts | 26 ++++++ test/session-manager-auto-recover.test.ts | 58 +++++++++--- 13 files changed, 369 insertions(+), 88 deletions(-) create mode 100644 test/dashboard-row-lightweight-source.test.ts diff --git a/src/config.ts b/src/config.ts index c0740082da..5de1dea854 100644 --- a/src/config.ts +++ b/src/config.ts @@ -184,12 +184,11 @@ export const config = { cliId: (process.env.CLI_ID ?? 'claude-code') as import('./adapters/cli/types.js').CliId, cliPathOverride: process.env.CLI_PATH, backendType: (process.env.BACKEND_TYPE ?? detectDefaultBackend()) as BackendType, - /** Auto-recovery throttle: on restart every surviving persistent-backend - * session is eagerly re-forked to re-attach its pane. With dozens of - * sessions per daemon (and many daemons on one box) firing them all at - * once spikes CPU/IO, so the re-fork is staggered: spawn `batchSize` - * workers, wait `delayMs`, repeat. Tune via BOTMUX_RECOVERY_FORK_BATCH / - * BOTMUX_RECOVERY_FORK_DELAY_MS. */ + /** Auto-recovery re-attach is opt-in: keeping restored persistent sessions + * lazy protects daemon startup paths like message listeners from thousands + * of worker re-forks. Tune via BOTMUX_RECOVERY_FORK_ENABLED, + * BOTMUX_RECOVERY_FORK_BATCH, and BOTMUX_RECOVERY_FORK_DELAY_MS. */ + recoveryForkEnabled: (process.env.BOTMUX_RECOVERY_FORK_ENABLED ?? 'false').toLowerCase() === 'true', recoveryForkBatchSize: Math.max(1, Number(process.env.BOTMUX_RECOVERY_FORK_BATCH) || 5), recoveryForkDelayMs: Math.max(0, Number(process.env.BOTMUX_RECOVERY_FORK_DELAY_MS ?? 250)), forwardFollowupWaitMs: resolveForwardFollowupWaitMs(), diff --git a/src/core/dashboard-ipc-server.ts b/src/core/dashboard-ipc-server.ts index 985cf6e3b5..d58a1dc2a9 100644 --- a/src/core/dashboard-ipc-server.ts +++ b/src/core/dashboard-ipc-server.ts @@ -967,16 +967,18 @@ export { composeRowFromActive, composeRowFromClosed, composeRowFromPersistedActi // holder. export function setBotName(name: string): void { setRowsBotName(name); } +const DASHBOARD_SNAPSHOT_ROW_OPTS = { lightweight: true } as const; + function composeDashboardSessionRows(): SessionRow[] { - const active = listActiveSessions().map((ds) => composeRowFromActive(ds)); + const active = listActiveSessions().map(ds => composeRowFromActive(ds, DASHBOARD_SNAPSHOT_ROW_OPTS)); const activeIds = new Set(active.map(row => row.sessionId)); const persisted = sessionStore.listSessions(); const unregisteredActive = persisted .filter(session => session.status === 'active' && !activeIds.has(session.sessionId)) - .map(composeRowFromPersistedActive); + .map(session => composeRowFromPersistedActive(session, DASHBOARD_SNAPSHOT_ROW_OPTS)); const closed = persisted .filter(session => session.status === 'closed' && !activeIds.has(session.sessionId)) - .map(composeRowFromClosed); + .map(session => composeRowFromClosed(session, DASHBOARD_SNAPSHOT_ROW_OPTS)); return [...active, ...unregisteredActive, ...closed]; } @@ -5593,7 +5595,7 @@ ipcRoute('GET', '/api/events', (_req, res) => { const activeIds = new Set(); for (const ds of listActiveSessions()) { activeIds.add(ds.session.sessionId); - res.write(`event: session.spawned\ndata: ${JSON.stringify({ session: composeRowFromActive(ds) })}\n\n`); + res.write(`event: session.spawned\ndata: ${JSON.stringify({ session: composeRowFromActive(ds, DASHBOARD_SNAPSHOT_ROW_OPTS) })}\n\n`); } // Persisted active rows may be intentionally absent from the runtime Map // after an inconclusive exact-backend teardown. Replay them as dormant @@ -5601,7 +5603,7 @@ ipcRoute('GET', '/api/events', (_req, res) => { // /api/sessions and never synthesize a closed row. for (const s of sessionStore.listSessions()) { if (s.status !== 'active' || activeIds.has(s.sessionId)) continue; - res.write(`event: session.spawned\ndata: ${JSON.stringify({ session: composeRowFromPersistedActive(s) })}\n\n`); + res.write(`event: session.spawned\ndata: ${JSON.stringify({ session: composeRowFromPersistedActive(s, DASHBOARD_SNAPSHOT_ROW_OPTS) })}\n\n`); } // Also replay sessions CLOSED during this run as `session.spawned` carrying a // closed row. The active-only replay above can't cover a restore-time zombie: @@ -5621,7 +5623,7 @@ ipcRoute('GET', '/api/events', (_req, res) => { if (s.status !== 'closed' || activeIds.has(s.sessionId)) continue; const closedMs = s.closedAt ? Date.parse(s.closedAt) : NaN; if (!Number.isFinite(closedMs) || closedMs < PROCESS_START_MS) continue; - res.write(`event: session.spawned\ndata: ${JSON.stringify({ session: composeRowFromClosed(s) })}\n\n`); + res.write(`event: session.spawned\ndata: ${JSON.stringify({ session: composeRowFromClosed(s, DASHBOARD_SNAPSHOT_ROW_OPTS) })}\n\n`); } } catch (err) { logger.warn(`[dashboard-ipc] /api/events snapshot replay failed: ${err}`); diff --git a/src/core/dashboard-rows.ts b/src/core/dashboard-rows.ts index 5132480e3b..a14041150e 100644 --- a/src/core/dashboard-rows.ts +++ b/src/core/dashboard-rows.ts @@ -142,6 +142,11 @@ let cachedBotName = ''; export function setBotName(name: string): void { cachedBotName = name; } export function getBotName(): string { return cachedBotName; } +export interface ComposeSessionRowOptions { + fresh?: boolean; + lightweight?: boolean; +} + function parseSessionTime(iso: string | undefined): number | undefined { if (!iso) return undefined; const ms = Date.parse(iso); @@ -212,10 +217,13 @@ function sessionRuntimeFields(s: Session): Pick void, + batchSize: number = RECOVERY_FORK_BATCH_SIZE, + delayMs: number = RECOVERY_FORK_DELAY_MS, + stillOwned: (ds: DaemonSession) => boolean = ds => ds.session.status === 'active', +): void { + void staggeredRecoveryFork(sessions, fork, batchSize, delayMs, stillOwned) + .catch(err => { + logger.error( + `[restore] background recovery re-attach failed: ` + + `${err instanceof Error ? err.message : String(err)}`, + ); + }); +} + export async function restoreActiveSessions( activeSessions: Map, quarantinedSessionIds: ReadonlySet = new Set(), @@ -2139,7 +2152,7 @@ export async function restoreActiveSessions( continue; } restoredByThisInvocation.push(ds); - announceSessionRow(ds); + announceSessionRow(ds, { lightweight: true }); forkAdoptWorker(ds, { restoredFromMetadata: true }); logger.info(`[${session.sessionId.substring(0, 8)}] Restored adopt session (target: ${adoptTargetLabel(adopted)}, scope: ${scope})`); continue; @@ -2245,7 +2258,7 @@ export async function restoreActiveSessions( restoredByThisInvocation.push(ds); // 重启后把待办池卡片重新广播给 dashboard,否则会从看板消失(#277 同款修复, // 我这条 queued 分支提前 continue 绕过了下面的 announceSessionRow,要自己补)。 - announceSessionRow(ds); + announceSessionRow(ds, { lightweight: true }); if (restoredPendingRepo) { try { await resumeRestoredPendingRepoSetup(ds, activeSessions); @@ -2474,7 +2487,7 @@ export async function restoreActiveSessions( continue; } restoredByThisInvocation.push(ds); - announceSessionRow(ds); + announceSessionRow(ds, { lightweight: true }); if (session.initialUserTurnPending) { // `hasHistory: true` above means "there may be a CLI process/transcript to @@ -2518,6 +2531,7 @@ export async function restoreActiveSessions( backendName: string; }> = []; const namesByBackend = new Map>(); + const skippedReattachByBackend = new Map(); for (const ds of restoredByThisInvocation) { // A later restore CAS awaited after this row was registered. During that // yield the user may have closed/resumed/replaced it; never carry the stale @@ -2541,7 +2555,10 @@ export async function restoreActiveSessions( } continue; } - if (!shouldAutoForkOnRestore(backendType)) continue; + if (!shouldAutoForkOnRestore(backendType)) { + skippedReattachByBackend.set(backendType, (skippedReattachByBackend.get(backendType) ?? 0) + 1); + continue; + } // Honour the worker-selected target (Herdr may own an agent inside a shared // host session) rather than assuming the deterministic whole-session name. const backendTarget = persistentBackendTargetForSession(ds)!; @@ -2556,6 +2573,16 @@ export async function restoreActiveSessions( names.add(backendTarget.sessionName); namesByBackend.set(backendType, names); } + if (skippedReattachByBackend.size > 0) { + const total = [...skippedReattachByBackend.values()].reduce((sum, count) => sum + count, 0); + const detail = [...skippedReattachByBackend.entries()] + .map(([backendType, count]) => `${backendType}=${count}`) + .join(', '); + logger.info( + `[restore] skipped eager re-attach for ${total} persistent session(s) ` + + `(${detail}); lazy recovery remains available`, + ); + } // ZMX/Zellij can classify every requested name from one control-plane list. // This is both a consistent restore snapshot and avoids an O(N²) ZMX restart // when each per-row probe would otherwise scan every per-session daemon. @@ -2658,11 +2685,7 @@ export async function restoreActiveSessions( toReattach.push(ds); } - // Staggered re-fork (see staggeredRecoveryFork): empty prompt = re-attach - // only, no new turn — same as the old per-session eager fork. - await staggeredRecoveryFork( - toReattach, - (ds) => { + const reattach = (ds: DaemonSession): void => { // A quarantined tail-only owner (restore promotion failed transiently) is // handled by the CENTRAL guard inside forkWorker: this blank fork retries // the old head's promotion first and, if it still fails, refuses to fork @@ -2686,11 +2709,24 @@ export async function restoreActiveSessions( } : true, ); - }, - RECOVERY_FORK_BATCH_SIZE, - RECOVERY_FORK_DELAY_MS, - ds => activeSessions.get(activeSessionKey(ds)) === ds, - ); + }; + + // Staggered re-fork (see staggeredRecoveryFork): empty prompt = re-attach + // only, no new turn — same as the old per-session eager fork. Keep it off the + // restore critical path: on long-lived installations thousands of restored + // active rows can otherwise hold daemon readiness and message-listener + // backfill behind worker re-attach for minutes, even though the sessions are + // already registered and can lazy cold-resume on demand. + if (toReattach.length > 0) { + logger.info(`[restore] scheduling ${toReattach.length} persistent session(s) for background re-attach`); + scheduleStaggeredRecoveryFork( + toReattach, + reattach, + RECOVERY_FORK_BATCH_SIZE, + RECOVERY_FORK_DELAY_MS, + ds => activeSessions.get(activeSessionKey(ds)) === ds, + ); + } const hasPersistentBackend = [...activeSessions.values()].some(ds => !!getSessionPersistentBackendType(ds)); logger.info(`Restored ${active.length} session(s)${hasPersistentBackend ? '' : ', waiting for messages to resume'}`); diff --git a/src/im/lark/client.ts b/src/im/lark/client.ts index a038e9406b..14b09018f4 100644 --- a/src/im/lark/client.ts +++ b/src/im/lark/client.ts @@ -1611,6 +1611,10 @@ export async function listChatMessages( export interface ChatMessageScanOptions { /** Lark page size per request. Clamped to the API max of 50. */ pageSize?: number; + /** Optional HTTP timeout for each page request. Falls back to the client default. */ + timeoutMs?: number; + /** Optional cancellation signal for the underlying page requests. */ + signal?: AbortSignal; /** * Called while scanning newest -> oldest. Returning true stops after the * current message has been included in the returned chronological list. @@ -1639,7 +1643,7 @@ export async function listChatMessagesUntil( sort_type: 'ByCreateTimeDesc', with_sender_name: 'true', ...(pageToken ? { page_token: pageToken } : {}), - }); + }, { timeoutMs: options.timeoutMs, signal: options.signal }); if (res.code !== 0) { throw new Error(`Failed to list chat messages: ${res.msg} (code: ${res.code})`); diff --git a/src/im/lark/event-dispatcher.ts b/src/im/lark/event-dispatcher.ts index fb3263cadf..327d4690f9 100644 --- a/src/im/lark/event-dispatcher.ts +++ b/src/im/lark/event-dispatcher.ts @@ -137,6 +137,16 @@ export function writeBotInfoFile(dataDir: string): void { /** Per-app in-flight open_id probe, so a startup burst of events shares one probe. */ const inflightOpenIdProbes = new Map>(); +async function fetchWithTimeout(input: RequestInfo | URL, init: RequestInit, timeoutMs: number): Promise { + const controller = new AbortController(); + const timer = setTimeout(() => controller.abort(), timeoutMs); + try { + return await fetch(input, { ...init, signal: controller.signal }); + } finally { + clearTimeout(timer); + } +} + /** * Ensure the bot's own open_id is resolved before @-detection. `probeBotOpenId` * is fired fire-and-forget at daemon startup, so events can arrive while @@ -163,19 +173,19 @@ export async function probeBotOpenId(larkAppId: string): Promise { const openApi = larkHosts(normalizeBrand(bot.config.brand)).openApi; // Call /bot/v3/info to get the bot's open_id using tenant_access_token - const tokenRes = await fetch(`${openApi}/open-apis/auth/v3/tenant_access_token/internal`, { + const tokenRes = await fetchWithTimeout(`${openApi}/open-apis/auth/v3/tenant_access_token/internal`, { method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ app_id: bot.config.larkAppId, app_secret: bot.config.larkAppSecret }), - }); + }, 10_000); const tokenData = await tokenRes.json() as any; if (tokenData.code !== 0) { throw new Error(`Failed to get tenant_access_token: ${tokenData.msg}`); } - const botRes = await fetch(`${openApi}/open-apis/bot/v3/info/`, { + const botRes = await fetchWithTimeout(`${openApi}/open-apis/bot/v3/info/`, { headers: { Authorization: `Bearer ${tokenData.tenant_access_token}` }, - }); + }, 10_000); const botData = await botRes.json() as any; if (botData.code !== 0) { throw new Error(`Failed to get bot info: ${botData.msg}`); @@ -2146,6 +2156,40 @@ const MESSAGE_LISTENER_BACKFILL_PAGE_SIZE = Math.min(50, Math.max( 1, Number(process.env.BOTMUX_MESSAGE_LISTENER_BACKFILL_PAGE_SIZE) || 50, )); +const MESSAGE_LISTENER_BACKFILL_TIMEOUT_MS = Math.max( + 1_000, + Number(process.env.BOTMUX_MESSAGE_LISTENER_BACKFILL_TIMEOUT_MS) || 15_000, +); +const MESSAGE_LISTENER_POLL_RUN_TIMEOUT_MS = Math.max( + MESSAGE_LISTENER_BACKFILL_TIMEOUT_MS + 5_000, + Number(process.env.BOTMUX_MESSAGE_LISTENER_POLL_RUN_TIMEOUT_MS) || 30_000, +); + +function withMessageListenerTimeout(promise: Promise, timeoutMs: number, label: string): Promise { + let timer: ReturnType | undefined; + return Promise.race([ + promise, + new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error(`${label} timed out after ${timeoutMs}ms`)), timeoutMs); + }), + ]).finally(() => { + if (timer) clearTimeout(timer); + }); +} + +async function withAbortableMessageListenerTimeout( + run: (signal: AbortSignal) => Promise, + timeoutMs: number, + label: string, +): Promise { + const controller = new AbortController(); + try { + return await withMessageListenerTimeout(run(controller.signal), timeoutMs, label); + } catch (err) { + controller.abort(); + throw err; + } +} function enabledMessageListenerChatIds(bot: BotState): string[] { return Object.entries(bot.config.messageListeners ?? {}) @@ -2198,10 +2242,14 @@ function larkReceiveEventFromHistoryMessage(message: any, chatId: string): any { const isBotSenderType = senderIdType === 'app_id' || senderTypeRaw === 'app' || senderTypeRaw === 'bot'; const isOpenIdDomain = senderIdType === 'open_id' || (typeof senderOpenId === 'string' && senderOpenId.startsWith('ou_')); + const content = typeof message?.content === 'string' + ? message.content + : typeof message?.body?.content === 'string' ? message.body.content : ''; return { message: { ...message, message_type: message?.message_type ?? message?.msg_type, + content, chat_id: message?.chat_id ?? chatId, chat_type: message?.chat_type ?? 'group', }, @@ -2264,19 +2312,26 @@ async function pollMessageListenersOnce(larkAppId: string, handlers: EventHandle if (chatIds.length === 0) return; const cutoff = now - MESSAGE_LISTENER_BACKFILL_WINDOW_MS; - await ensureBotOpenId(larkAppId).catch(() => { /* degrade; heartbeat retries */ }); + await withMessageListenerTimeout(ensureBotOpenId(larkAppId), 10_000, 'ensure bot open_id') + .catch(err => logger.warn(`[message-listener:${larkAppId}] bot open_id probe skipped: ${err instanceof Error ? err.message : String(err)}`)); for (const chatId of chatIds) { let messages: any[]; try { - messages = await listChatMessagesUntil(larkAppId, chatId, { - pageSize: MESSAGE_LISTENER_BACKFILL_PAGE_SIZE, - stopAfter: (message, seenCount) => { - const createdAt = messageCreateTimeMs(message); - return seenCount >= MESSAGE_LISTENER_BACKFILL_SCAN_LIMIT || - (Number.isFinite(createdAt) && (createdAt as number) < cutoff); - }, - }); + messages = await withAbortableMessageListenerTimeout( + signal => listChatMessagesUntil(larkAppId, chatId, { + pageSize: MESSAGE_LISTENER_BACKFILL_PAGE_SIZE, + timeoutMs: MESSAGE_LISTENER_BACKFILL_TIMEOUT_MS, + signal, + stopAfter: (message, seenCount) => { + const createdAt = messageCreateTimeMs(message); + return seenCount >= MESSAGE_LISTENER_BACKFILL_SCAN_LIMIT || + (Number.isFinite(createdAt) && (createdAt as number) < cutoff); + }, + }), + MESSAGE_LISTENER_BACKFILL_TIMEOUT_MS, + `poll chat ${chatId.substring(0, 12)}`, + ); } catch (err) { logger.warn(`[message-listener:${larkAppId}] failed to poll chat=${chatId.substring(0, 12)}: ${err instanceof Error ? err.message : String(err)}`); continue; @@ -2324,6 +2379,14 @@ export async function __pollMessageListenersOnceForTest(larkAppId: string, handl await pollMessageListenersOnce(larkAppId, handlers, now); } +function runMessageListenerPoll(larkAppId: string, handlers: EventHandlers, label: string): Promise { + return withMessageListenerTimeout( + pollMessageListenersOnce(larkAppId, handlers), + MESSAGE_LISTENER_POLL_RUN_TIMEOUT_MS, + `${label} message-listener poll`, + ); +} + function usesForwardFollowupDelay(mentionMode: GroupMentionMode): boolean { return mentionMode === 'never' || mentionMode === 'ambient'; } @@ -4174,7 +4237,7 @@ export function startLarkEventDispatcher(larkAppId: string, larkAppSecret: strin const listenerPollTimer = setInterval(() => { if (listenerPollInFlight) return; listenerPollInFlight = true; - void pollMessageListenersOnce(larkAppId, handlers) + void runMessageListenerPoll(larkAppId, handlers, 'scheduled') .catch(err => logger.error(`[message-listener:${larkAppId}] poll failed: ${err instanceof Error ? err.message : String(err)}`)) .finally(() => { listenerPollInFlight = false; }); }, MESSAGE_LISTENER_POLL_INTERVAL_MS); @@ -4184,7 +4247,7 @@ export function startLarkEventDispatcher(larkAppId: string, larkAppSecret: strin setTimeout(() => { if (listenerPollInFlight) return; listenerPollInFlight = true; - void pollMessageListenersOnce(larkAppId, handlers) + void runMessageListenerPoll(larkAppId, handlers, 'initial') .catch(err => logger.error(`[message-listener:${larkAppId}] initial poll failed: ${err instanceof Error ? err.message : String(err)}`)) .finally(() => { listenerPollInFlight = false; }); }, 2_000).unref(); diff --git a/src/im/lark/message-parser.ts b/src/im/lark/message-parser.ts index 6450a4984e..d5eedfd3fc 100644 --- a/src/im/lark/message-parser.ts +++ b/src/im/lark/message-parser.ts @@ -25,7 +25,8 @@ interface RawEventData { thread_id?: string; parent_id?: string; message_type: string; // NOT msg_type - content: string; + content?: string; + body?: { content?: string }; chat_id: string; chat_type: string; create_time: string; @@ -123,6 +124,11 @@ export function mentionIdType(m: any): string | undefined { return undefined; } +function rawEventMessageContent(message: RawEventData['message']): string { + const raw = message.content ?? message.body?.content; + return typeof raw === 'string' ? raw : ''; +} + /** * Extract a bot mention's app_id across ALL shapes Lark has been observed to * use. app_id-form bot mentions must never flow through mentionOpenId() (which @@ -335,7 +341,7 @@ export function unwrapUserDslContent(rawContent: string): string | null { } function unwrapUserDsl(data: RawEventData): boolean { - const unwrapped = unwrapUserDslContent(data.message.content); + const unwrapped = unwrapUserDslContent(rawEventMessageContent(data.message)); if (unwrapped === null) return false; data.message.content = unwrapped; logger.info(`[parser] Unwrapped user_dsl for ${data.message.message_id}`); @@ -497,12 +503,14 @@ export function parseEventMessage( numberer: ImgNumberer = createImgNumberer(), ): { parsed: LarkMessage; resources: MessageResource[] } { const { sender, message } = data; + const rawContent = rawEventMessageContent(message); + const normalizedContent = normalizeApiMessageContent(message.message_type, rawContent); // Trace non-text messages at debug only (DEBUG=1 gated). The raw card/post // content can be many KB and may include attachment metadata that's // noisy or sensitive — truncate to ~500 chars so DEBUG logs stay scannable. if (message.message_type !== 'text' && logger.isDebug()) { - const raw = message.content ?? ''; + const raw = rawContent; const trimmed = raw.length > 500 ? raw.slice(0, 500) + `…(+${raw.length - 500}b)` : raw; logger.debug(`[parser] type=${message.message_type} content=${trimmed} keys=${Object.keys(message).join(',')}`); } @@ -510,7 +518,7 @@ export function parseEventMessage( // Share numberer so in-body [图片 N] placeholders use the same numbers as // the attachment list. Resources first → numbers assigned; text second → // reuses them. - const resources = extractResources(message.message_type, message.content, numberer); + const resources = extractResources(message.message_type, rawContent, numberer); // Extract structured mentions const mentions: LarkMention[] | undefined = @@ -535,7 +543,7 @@ export function parseEventMessage( senderUnionId: sender.sender_id?.union_id, senderType: sender.sender_type, msgType: message.message_type, - content: extractTextContent(message.message_type, message.content, message.mentions, numberer), + content: extractTextContent(message.message_type, normalizedContent, message.mentions, numberer), createTime: message.create_time, mentions, }; diff --git a/test/dashboard-attention-signals.test.ts b/test/dashboard-attention-signals.test.ts index ef70694635..c44f0c140f 100644 --- a/test/dashboard-attention-signals.test.ts +++ b/test/dashboard-attention-signals.test.ts @@ -125,6 +125,14 @@ describe('attention signals', () => { expect(composeRowFromActive(queued).status).toBe('idle'); }); + it('composeRowFromActive can skip restore-time expensive enrichment', () => { + const row = composeRowFromActive(makeDs(), { lightweight: true }); + + expect(row.tokenUsage).toBeUndefined(); + expect(row.previewUserText).toBeUndefined(); + expect(row.previewBotText).toBeUndefined(); + }); + it('composeRowFromActive carries the agent raise-hand signal with its reason', () => { const raised = composeRowFromActive(makeDs({ agentAttention: { kind: 'authz', reason: '需要 prod 部署授权', at: 1234 }, diff --git a/test/dashboard-row-lightweight-source.test.ts b/test/dashboard-row-lightweight-source.test.ts new file mode 100644 index 0000000000..4014f1573d --- /dev/null +++ b/test/dashboard-row-lightweight-source.test.ts @@ -0,0 +1,33 @@ +import { readFileSync } from 'node:fs'; +import { resolve } from 'node:path'; +import { describe, expect, it } from 'vitest'; + +const source = readFileSync(resolve(process.cwd(), 'src/core/dashboard-ipc-server.ts'), 'utf8'); +const rowsSource = readFileSync(resolve(process.cwd(), 'src/core/dashboard-rows.ts'), 'utf8'); + +describe('dashboard bulk session snapshots', () => { + it('uses lightweight rows for /api/sessions snapshots', () => { + expect(source).toContain('const DASHBOARD_SNAPSHOT_ROW_OPTS = { lightweight: true } as const;'); + expect(source).toContain('composeRowFromActive(ds, DASHBOARD_SNAPSHOT_ROW_OPTS)'); + expect(source).toContain('composeRowFromPersistedActive(session, DASHBOARD_SNAPSHOT_ROW_OPTS)'); + expect(source).toContain('composeRowFromClosed(session, DASHBOARD_SNAPSHOT_ROW_OPTS)'); + }); + + it('uses lightweight rows for /api/events replay snapshots', () => { + expect(source).toContain('composeRowFromActive(ds, DASHBOARD_SNAPSHOT_ROW_OPTS) })}\\n\\n`);'); + expect(source).toContain('composeRowFromPersistedActive(s, DASHBOARD_SNAPSHOT_ROW_OPTS) })}\\n\\n`);'); + expect(source).toContain('composeRowFromClosed(s, DASHBOARD_SNAPSHOT_ROW_OPTS) })}\\n\\n`);'); + }); + + it('keeps open todo extraction behind the lightweight guard', () => { + const start = rowsSource.indexOf('export function composeRowFromActive'); + const end = rowsSource.indexOf('export function composeRowFromClosed'); + const activeComposer = rowsSource.slice(start, end); + const guard = 'if (!opts.lightweight) {'; + const beforeGuard = activeComposer.slice(0, activeComposer.indexOf(guard)); + + expect(activeComposer).toContain(guard); + expect(activeComposer).toContain('row.openTodos = sessionOpenTodos(ds.session, ds.workingDir, opts.fresh);'); + expect(beforeGuard).not.toContain('sessionOpenTodos('); + }); +}); diff --git a/test/event-dispatcher.test.ts b/test/event-dispatcher.test.ts index 028f74f5cd..66bbc1185f 100644 --- a/test/event-dispatcher.test.ts +++ b/test/event-dispatcher.test.ts @@ -1606,7 +1606,7 @@ describe('message listener polling backfill', () => { await flushEventWork(); expect(mockListChatMessagesUntil).toHaveBeenCalledWith(MY_APP_ID, 'chat_listener', expect.objectContaining({ - pageSize: expect.any(Number), + pageSize: 50, stopAfter: expect.any(Function), })); expect(handlers.handleNewTopic).toHaveBeenCalledWith( @@ -1625,6 +1625,30 @@ describe('message listener polling backfill', () => { ); }); + it('passes an abort signal to polled chat-history requests', async () => { + const card = makeHistoryMessage({ + senderAppId: OTHER_BOT_APP_ID, + senderType: 'app', + messageType: 'interactive', + messageId: 'msg-polled-abort-signal', + chatId: 'chat_listener', + content: JSON.stringify({ title: 'Argos平台报警', elements: [[{ tag: 'text', text: '需要可取消的轮询' }]] }), + createTime: String(Date.now()), + }); + let signal: AbortSignal | undefined; + mockListChatMessagesUntil.mockImplementationOnce(async (_appId, _chatId, options) => { + signal = options?.signal; + return [card]; + }); + + await __pollMessageListenersOnceForTest(MY_APP_ID, handlers); + await flushEventWork(); + + expect(signal).toBeDefined(); + expect(typeof signal?.addEventListener).toBe('function'); + expect(signal?.aborted).toBe(false); + }); + it('does not replay a polled listener message after the message_id is claimed', async () => { const card = makeHistoryMessage({ senderAppId: OTHER_BOT_APP_ID, @@ -1644,6 +1668,35 @@ describe('message listener polling backfill', () => { expect(handlers.handleNewTopic).toHaveBeenCalledTimes(1); }); + it('normalizes REST body.content before dispatching a polled listener match', async () => { + const cardContent = JSON.stringify({ + title: 'Argos平台报警', + elements: [[{ tag: 'text', text: 'body-only 告警详情' }]], + }); + const card = { + message_id: 'msg-polled-body-only', + chat_id: 'chat_listener', + chat_type: 'group', + msg_type: 'interactive', + body: { content: cardContent }, + create_time: String(Date.now()), + sender: { + id: OTHER_BOT_APP_ID, + id_type: 'app_id', + sender_type: 'app', + }, + }; + mockListChatMessagesUntil.mockResolvedValueOnce([card]); + handlers.handleNewTopic.mockImplementationOnce(async (data) => { + expect((data as any).message.content).toBe(cardContent); + }); + + await __pollMessageListenersOnceForTest(MY_APP_ID, handlers); + await flushEventWork(); + + expect(handlers.handleNewTopic).toHaveBeenCalledTimes(1); + }); + it('resolves a sibling bot app_id to open_id so an open_id include list matches on the polled path', async () => { // Realistic config: the sender filter stores the peer bot's OPEN_ID (what the // dashboard member picker saves), while chat history reports the bot by app_id. diff --git a/test/message-parser.test.ts b/test/message-parser.test.ts index dee5d1c6d4..2230ba594c 100644 --- a/test/message-parser.test.ts +++ b/test/message-parser.test.ts @@ -59,6 +59,32 @@ describe('parseApiMessage metadata', () => { }); }); +describe('parseEventMessage compatibility', () => { + it('parses REST history messages that carry content in body.content', () => { + const result = parseEventMessage({ + sender: { + sender_type: 'app', + sender_id: { app_id: 'cli_argos' }, + }, + message: { + message_id: 'om_body_only', + chat_id: 'oc_alerts', + chat_type: 'group', + message_type: 'interactive', + create_time: '1000', + body: { + content: JSON.stringify({ + title: 'Argos平台报警', + elements: [[{ tag: 'text', text: 'body-only 告警详情' }]], + }), + }, + }, + } as any); + + expect(result.parsed.content).toContain('body-only 告警详情'); + }); +}); + // ─── Interactive card: Format A (Lark API simplified) ───────────────────── describe('Interactive card parsing: Format A (API simplified)', () => { diff --git a/test/session-manager-auto-recover.test.ts b/test/session-manager-auto-recover.test.ts index 055f4d6022..7f865a2893 100644 --- a/test/session-manager-auto-recover.test.ts +++ b/test/session-manager-auto-recover.test.ts @@ -1,12 +1,10 @@ /** * Auto-recovery on daemon restart. * - * On restart every surviving persistent-backend session is eagerly re-forked to - * re-attach its pane, so the session actually comes back instead of sitting dead - * until its next message (and a pane whose CLI died gets healed, keeping the - * transcript fallback working). The old `BOTMUX_QUIET_RESTART` gate that - * suppressed this is gone — card silence is now handled by `suppressRecoveryCard` - * on restored sessions, not by skipping recovery. + * On restart surviving persistent-backend sessions are kept lazy by default: + * restoring thousands of active rows must not starve message listeners by + * eagerly re-forking every worker. Operators can explicitly opt in to eager + * re-attach when terminal readiness is more important than startup latency. * * `staggeredRecoveryFork` spaces the re-forks out (batch + delay) so a box with * dozens of surviving sessions doesn't spike on restart, and skips any session @@ -21,24 +19,38 @@ vi.mock('../src/bot-registry.js', () => ({ vi.mock('../src/config.js', () => ({ config: { - daemon: { workingDir: '~', workingDirs: ['~'], recoveryForkBatchSize: 5, recoveryForkDelayMs: 0 }, + daemon: { + workingDir: '~', + workingDirs: ['~'], + recoveryForkBatchSize: 5, + recoveryForkDelayMs: 0, + recoveryForkEnabled: false, + }, session: { dataDir: '/tmp/botmux-test' }, }, })); -import { shouldAutoForkOnRestore, staggeredRecoveryFork } from '../src/core/session-manager.js'; +import { shouldAutoForkOnRestore, staggeredRecoveryFork, scheduleStaggeredRecoveryFork } from '../src/core/session-manager.js'; import type { DaemonSession } from '../src/core/types.js'; describe('shouldAutoForkOnRestore', () => { - it('eagerly re-forks every persistent backend (tmux/herdr/zellij/zmx)', () => { - expect(shouldAutoForkOnRestore('tmux')).toBe(true); - expect(shouldAutoForkOnRestore('herdr')).toBe(true); - expect(shouldAutoForkOnRestore('zellij')).toBe(true); - expect(shouldAutoForkOnRestore('zmx')).toBe(true); + it('keeps persistent backends lazy by default so startup does not starve listeners', () => { + expect(shouldAutoForkOnRestore('tmux')).toBe(false); + expect(shouldAutoForkOnRestore('herdr')).toBe(false); + expect(shouldAutoForkOnRestore('zellij')).toBe(false); + expect(shouldAutoForkOnRestore('zmx')).toBe(false); + }); + + it('can explicitly eager re-attach persistent backends', () => { + expect(shouldAutoForkOnRestore('tmux', true)).toBe(true); + expect(shouldAutoForkOnRestore('herdr', true)).toBe(true); + expect(shouldAutoForkOnRestore('zellij', true)).toBe(true); + expect(shouldAutoForkOnRestore('zmx', true)).toBe(true); }); it('never eagerly forks the pty backend — it has no pane to re-attach', () => { expect(shouldAutoForkOnRestore('pty')).toBe(false); + expect(shouldAutoForkOnRestore('pty', true)).toBe(false); }); }); @@ -112,4 +124,24 @@ describe('staggeredRecoveryFork', () => { expect(forked).toEqual(['a']); }); + + it('can be scheduled without blocking restore on delayed batches', async () => { + vi.useFakeTimers(); + try { + const sessions = Array.from({ length: 5 }, (_, i) => ds(`s${i}`)); + const forked: string[] = []; + + scheduleStaggeredRecoveryFork(sessions, (d) => forked.push(d.session.sessionId), 2, 20); + + expect(forked).toEqual(['s0', 's1']); + await vi.advanceTimersByTimeAsync(19); + expect(forked).toEqual(['s0', 's1']); + await vi.advanceTimersByTimeAsync(1); + expect(forked).toEqual(['s0', 's1', 's2', 's3']); + await vi.advanceTimersByTimeAsync(20); + expect(forked).toEqual(['s0', 's1', 's2', 's3', 's4']); + } finally { + vi.useRealTimers(); + } + }); }); From f65839a62a711aff91bbf041ecc09ab1eb9c79ba Mon Sep 17 00:00:00 2001 From: "zhujinhao.zjh" Date: Wed, 26 Aug 2026 14:43:45 +0800 Subject: [PATCH 2/2] Add listener session retention cleanup --- src/bot-registry.ts | 9 ++ src/core/dashboard-ipc-server.ts | 1 + src/daemon.ts | 80 ++++++++++++- src/dashboard/web/i18n.ts | 4 + src/dashboard/web/roles-page.tsx | 59 +++++++++ src/dashboard/web/roles.ts | 12 +- .../message-listener-session-cleanup.ts | 73 ++++++++++++ src/services/message-listener-store.ts | 14 ++- src/services/session-store.ts | 33 ++++++ src/types.ts | 4 + test/bot-registry.test.ts | 1 + test/message-listener-session-cleanup.test.ts | 112 ++++++++++++++++++ test/message-listener-store.test.ts | 57 ++++++++- test/session-store.test.ts | 30 +++++ 14 files changed, 484 insertions(+), 5 deletions(-) create mode 100644 src/services/message-listener-session-cleanup.ts create mode 100644 test/message-listener-session-cleanup.test.ts diff --git a/src/bot-registry.ts b/src/bot-registry.ts index 2d4cb7f8b0..81e50677db 100644 --- a/src/bot-registry.ts +++ b/src/bot-registry.ts @@ -39,6 +39,7 @@ import { normalizeSessionOwnerReminderConfig, type SessionOwnerReminderConfig, } from './core/session-owner-reminder.js'; +import { normalizeMessageListenerCleanupConfig } from './services/message-listener-session-cleanup.js'; import type { VcMeetingConsumerAgentConfig, VcMeetingConsumerConfig, @@ -237,6 +238,12 @@ export interface MessageListenerConfig { /** V1 starts one session per matched message. */ sessionMode?: 'per_message'; }; + cleanup?: { + /** Default true. */ + enabled?: boolean; + /** Default 168 hours. */ + retentionHours?: number; + }; } export interface SummaryRangeConfig { @@ -1060,6 +1067,7 @@ function normalizeMessageListenerConfig(raw: unknown, botIndex: number, chatId: const includeMsgTypes = normalizeMessageListenerStringList(messageRaw.includeMsgTypes); if (includeMsgTypes) messagePolicy.includeMsgTypes = includeMsgTypes; messagePolicy.scope = 'top_level'; + const cleanup = normalizeMessageListenerCleanupConfig(entry.cleanup); const contentRaw = entry.contentPolicy && typeof entry.contentPolicy === 'object' && !Array.isArray(entry.contentPolicy) ? entry.contentPolicy as Record @@ -1089,6 +1097,7 @@ function normalizeMessageListenerConfig(raw: unknown, botIndex: number, chatId: ...(Object.keys(messagePolicy).length > 0 ? { messagePolicy } : {}), ...(contentPolicy ? { contentPolicy } : {}), replyPolicy: { mode: 'thread', sessionMode: 'per_message' }, + cleanup, }; } diff --git a/src/core/dashboard-ipc-server.ts b/src/core/dashboard-ipc-server.ts index d58a1dc2a9..80f9bd4f81 100644 --- a/src/core/dashboard-ipc-server.ts +++ b/src/core/dashboard-ipc-server.ts @@ -3521,6 +3521,7 @@ async function collectMessageListenerPreviewMatches( 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' } }), + ...(listener.cleanup ? { cleanup: listener.cleanup } : {}), replyPolicy: { mode: 'thread', sessionMode: 'per_message' }, }; const previewBot = { diff --git a/src/daemon.ts b/src/daemon.ts index 412f674170..aa30b4bcea 100644 --- a/src/daemon.ts +++ b/src/daemon.ts @@ -192,6 +192,7 @@ import { readableTerminalUrlFor, findActiveBySessionId, withActiveSessionKeyLock, + destroyUnregisteredPersistentBacking, getDaemonBootId, getDaemonStreamingCardUsageSnapshot, postTurnStartingCard, @@ -283,6 +284,10 @@ import { } from './core/session-title.js'; import { settleDeferredScheduleRun } from './core/deferred-schedule-settlement.js'; import { renderMessageListenerPrompt, refreshListenerCardTextFromResolved } from './services/message-listener.js'; +import { + MESSAGE_LISTENER_CLEANUP_INTERVAL_MS, + selectExpiredMessageListenerSessions, +} from './services/message-listener-session-cleanup.js'; import { sweepOrphanSandboxes } from './adapters/backend/sandbox.js'; import { TmuxBackend } from './adapters/backend/tmux-backend.js'; import { HerdrBackend } from './adapters/backend/herdr-backend.js'; @@ -3359,6 +3364,74 @@ function startMemoryDiagnostics(): ReturnType | undefined { return timer; } +function createMessageListenerSessionCleanupRunner(larkAppId: string): { + run(reason: 'startup' | 'interval'): Promise; + start(): ReturnType; +} { + let running = false; + const run = async (reason: 'startup' | 'interval'): Promise => { + if (running) return; + running = true; + try { + const bot = getBot(larkAppId); + const expired = selectExpiredMessageListenerSessions({ + sessions: sessionStore.listSessions(), + listeners: bot.config.messageListeners, + }); + if (expired.length === 0) return; + let closed = 0; + let failed = 0; + const inactive: string[] = []; + logger.info(`[message-listener-cleanup] ${reason}: closing ${expired.length} expired listener session(s) for ${larkAppId}`); + for (const session of expired) { + if (!findActiveBySessionId(session.sessionId)) { + try { + destroyUnregisteredPersistentBacking(session); + inactive.push(session.sessionId); + } catch (err) { + failed += 1; + logger.warn( + `[message-listener-cleanup] backing cleanup failed session=${session.sessionId}: ` + + `${err instanceof Error ? err.message : String(err)}`, + ); + } + continue; + } + try { + const result = await closeSessionHelper(session.sessionId); + if (result.ok && result.known) closed += 1; + else failed += 1; + } catch (err) { + failed += 1; + logger.warn( + `[message-listener-cleanup] close failed session=${session.sessionId}: ` + + `${err instanceof Error ? err.message : String(err)}`, + ); + } + } + if (inactive.length > 0) { + const inactiveIds = new Set(inactive); + closed += sessionStore.closeSessionsMatching(session => inactiveIds.has(session.sessionId)); + } + logger.info(`[message-listener-cleanup] ${reason}: closed=${closed} failed=${failed} bot=${larkAppId}`); + } catch (err) { + logger.warn(`[message-listener-cleanup] ${reason} failed: ${err instanceof Error ? err.message : String(err)}`); + } finally { + running = false; + } + }; + return { + run, + start() { + const timer = setInterval(() => { + void run('interval'); + }, MESSAGE_LISTENER_CLEANUP_INTERVAL_MS); + timer.unref?.(); + return timer; + }, + }; +} + /** * Reply into a session — scope-aware. * @@ -18199,6 +18272,7 @@ async function handleNewTopicAdmitted(data: any, ctx: RoutingContext): Promise { sessionStore.init(cfg.larkAppId); chatFirstSeenStore.init(cfg.larkAppId); initSessionGroups(cfg.larkAppId); + const messageListenerCleanup = createMessageListenerSessionCleanupRunner(cfg.larkAppId); const ambiguousOnBoot = reconcileVcMeetingDeliveriesOnBoot( config.session.dataDir, { receiverBootId: getDaemonBootId(), agentAppId: cfg.larkAppId }, @@ -22307,7 +22382,7 @@ export async function startDaemon(botIndex?: number): Promise { reapOrphanWorkers(); // Restore active sessions from previous run - // Restore active sessions from previous run + await messageListenerCleanup.run('startup'); await restoreActiveSessions(activeSessions, idempotencyQuarantinedSessionIds); // Restore complete → /api/asks may now safely 403 unknown sessions again; a // reconnecting ask hook that raced the restore got retryable 503s until here. @@ -22465,6 +22540,7 @@ export async function startDaemon(botIndex?: number): Promise { enforceLiveSessionCap('periodic'); }, 60_000); idleWorkerSweepTimer.unref?.(); + const messageListenerCleanupTimer = messageListenerCleanup.start(); const sessionOwnerReminder = new SessionOwnerReminderController({ load: () => loadSessionOwnerReminderRecords(config.session.dataDir, cfg.larkAppId), @@ -22862,6 +22938,7 @@ export async function startDaemon(botIndex?: number): Promise { clearInterval(descriptorHeartbeat); clearInterval(idleWorkerSweepTimer); clearInterval(sessionOwnerReminderTimer); + clearInterval(messageListenerCleanupTimer); if (memoryDiagnostics) clearInterval(memoryDiagnostics); removeDaemonDescriptor(cfg.larkAppId); ipcHandle.close().catch(() => { /* swallow */ }); @@ -23074,6 +23151,7 @@ export async function startDaemon(botIndex?: number): Promise { clearInterval(descriptorHeartbeat); clearInterval(idleWorkerSweepTimer); clearInterval(sessionOwnerReminderTimer); + clearInterval(messageListenerCleanupTimer); clearInterval(docCommentPollTimer); if (memoryDiagnostics) clearInterval(memoryDiagnostics); removeDaemonDescriptor(cfg.larkAppId); diff --git a/src/dashboard/web/i18n.ts b/src/dashboard/web/i18n.ts index 677a972d35..eaccc6c7f9 100644 --- a/src/dashboard/web/i18n.ts +++ b/src/dashboard/web/i18n.ts @@ -2435,6 +2435,8 @@ const zh = { 'roles.listenerReplyCardTitlePlaceholder': '留空使用默认标题', 'roles.listenerWorkingDir': '工作目录', 'roles.listenerWorkingDirPlaceholder': '留空使用 Bot 默认工作目录', + 'roles.listenerCleanupEnabled': '自动清理监听会话', + 'roles.listenerCleanupRetentionHours': '保留小时数', 'roles.listenerSenderTypes': '发送者类型', 'roles.listenerSenderUser': '用户', 'roles.listenerSenderBot': '机器人', @@ -4991,6 +4993,8 @@ const en: Record = { 'roles.listenerReplyCardTitlePlaceholder': 'Leave empty to use the default title', 'roles.listenerWorkingDir': 'Working directory', 'roles.listenerWorkingDirPlaceholder': 'Leave empty to use the bot default', + 'roles.listenerCleanupEnabled': 'Auto-clean listener sessions', + 'roles.listenerCleanupRetentionHours': 'Retention hours', 'roles.listenerSenderTypes': 'Sender types', 'roles.listenerSenderUser': 'Users', 'roles.listenerSenderBot': 'Bots', diff --git a/src/dashboard/web/roles-page.tsx b/src/dashboard/web/roles-page.tsx index c43a710025..b3ab348854 100644 --- a/src/dashboard/web/roles-page.tsx +++ b/src/dashboard/web/roles-page.tsx @@ -47,6 +47,7 @@ import { MAX_MESSAGE_LISTENER_PROMPT_BYTES, MAX_ROLE_BYTES, MESSAGE_LISTENER_WARN_BYTES, + DEFAULT_MESSAGE_LISTENER_CLEANUP_RETENTION_HOURS, DEFAULT_MESSAGE_LISTENER_PREVIEW_LIMIT, MAX_MESSAGE_LISTENER_PREVIEW_LIMIT, roleKey, @@ -116,6 +117,10 @@ const DEFAULT_LISTENER: MessageListenerData = { includeMsgTypes: [...LISTENER_MESSAGE_TYPES], scope: 'top_level', }, + cleanup: { + enabled: true, + retentionHours: DEFAULT_MESSAGE_LISTENER_CLEANUP_RETENTION_HOURS, + }, }; function cloneListener(listener: MessageListenerData | null | undefined): MessageListenerData { @@ -152,9 +157,24 @@ function cloneListener(listener: MessageListenerData | null | undefined): Messag ...(listener.contentPolicy.matchMode ? { matchMode: listener.contentPolicy.matchMode } : {}), }, } : {}), + cleanup: { + enabled: listener?.cleanup?.enabled !== false, + retentionHours: normalizeCleanupRetentionHours(listener?.cleanup?.retentionHours), + }, }; } +function normalizeCleanupRetentionHours(raw: unknown): number { + const value = typeof raw === 'string' && raw.trim() + ? Number(raw.trim()) + : raw; + if (typeof value !== 'number' || !Number.isFinite(value)) { + return DEFAULT_MESSAGE_LISTENER_CLEANUP_RETENTION_HOURS; + } + const normalized = Math.floor(value); + return normalized > 0 ? normalized : DEFAULT_MESSAGE_LISTENER_CLEANUP_RETENTION_HOURS; +} + function listenerHasConfig(listener: MessageListenerData | null): boolean { // A persisted listener is worth loading into the editor whenever it carries a // prompt — INCLUDING a disabled draft (enabled:false + non-empty prompt). The @@ -599,6 +619,17 @@ function RolesPage(props: { tab: RolesTab }) { })); } + function updateListenerCleanup(patch: Partial>): void { + setEditingListener(prev => ({ + ...prev, + cleanup: { + enabled: prev.cleanup?.enabled !== false, + retentionHours: normalizeCleanupRetentionHours(prev.cleanup?.retentionHours), + ...patch, + }, + })); + } + function toggleListenerSenderType(type: SenderTypeOption, checked: boolean): void { setEditingListener(prev => { const current = new Set(prev.senderPolicy?.includeSenderTypes ?? []); @@ -720,6 +751,10 @@ function RolesPage(props: { tab: RolesTab }) { ...(raw.matchMode === 'all' ? { matchMode: 'all' as const } : {}), }; })(); + const cleanup = { + enabled: editingListener.cleanup?.enabled !== false, + retentionHours: normalizeCleanupRetentionHours(editingListener.cleanup?.retentionHours), + }; return { enabled: editingListener.enabled, ...(editingListener.name?.trim() ? { name: editingListener.name.trim() } : {}), @@ -741,6 +776,7 @@ function RolesPage(props: { tab: RolesTab }) { scope: 'top_level', }, ...(contentPolicy ? { contentPolicy } : {}), + cleanup, }; } @@ -1203,6 +1239,7 @@ function RolesPage(props: { tab: RolesTab }) { onSenderPolicyPatch={updateListenerSenderPolicy} onMessagePolicyPatch={updateListenerMessagePolicy} onContentPolicyPatch={updateListenerContentPolicy} + onCleanupPatch={updateListenerCleanup} onToggleSenderType={toggleListenerSenderType} onToggleMsgType={toggleListenerMsgType} onSetTargetPolicy={setListenerTargetPolicy} @@ -1545,6 +1582,7 @@ function MessageListenerEditor(props: { onSenderPolicyPatch(patch: NonNullable): void; onMessagePolicyPatch(patch: NonNullable): void; onContentPolicyPatch(patch: NonNullable): void; + onCleanupPatch(patch: Partial>): void; onToggleSenderType(type: SenderTypeOption, checked: boolean): void; onToggleMsgType(msgType: string, checked: boolean): void; onSetTargetPolicy(openId: string, listening: boolean): void; @@ -1754,6 +1792,27 @@ function MessageListenerEditor(props: { +
+ + +
{tr('roles.listenerContentPolicy')}