diff --git a/resources/wcdb/win32/x64/wcdb_api.dll b/resources/wcdb/win32/x64/wcdb_api.dll index f1f9815..127dcc6 100644 Binary files a/resources/wcdb/win32/x64/wcdb_api.dll and b/resources/wcdb/win32/x64/wcdb_api.dll differ diff --git a/src/main/index.ts b/src/main/index.ts index b3d2ece..7530508 100644 --- a/src/main/index.ts +++ b/src/main/index.ts @@ -118,7 +118,7 @@ import { automationExecutionLogService } from './services/automation-execution-l import { initAutomationService, getAutomationService } from './services/automation-service' import { GroupStatsService } from './services/group-stats-service' import { wechatActionLogService } from './services/wechat-action-log-service' -import { wechatActionGateway } from './services/wechat-action-gateway' +import { toPersonalWechatSendResult, wechatActionGateway } from './services/wechat-action-gateway' import { personalWechatSendService } from './services/personal-wechat-send-service' import { getPersonalWechatSendCapability } from './services/personal-wechat-capability-service' import { scheduledReportService } from './services/scheduled-report-service' @@ -1056,14 +1056,17 @@ app.whenReady().then(async () => { voiceRecognition?.connect(voiceService, resolvedRoot) stickerService = new StickerService(wcdb4Client) videoAssetService = new VideoAssetService(wcdb4Client) - const monitoring = await wcdb4Client.startMonitor((type, json) => { + const monitoring = await wcdb4Client.startMonitor((event) => { wcdb4Client.invalidateSessionCache() // 正式 MessageListener:只做 coalesce + 有界回读 + dedup + 投递,不触发业务。 - messageListener.handleNativeChange() - groupExitMonitorService.notifyDatabaseChanged(json) - recallArchiveMonitor?.handleDatabaseChange(json) + // v2 事件带 sessionId ⇒ 回读不再依赖「最近活跃会话」。 + messageListener.handleNativeChange(event) + groupExitMonitorService.notifyDatabaseChanged(event.raw) + recallArchiveMonitor?.handleDatabaseChange(event.raw) for (const window of BrowserWindow.getAllWindows()) { - if (!window.isDestroyed()) window.webContents.send('wcdb-change', { type, json }) + if (!window.isDestroyed()) { + window.webContents.send('wcdb-change', { type: event.type, json: event.raw }) + } } }) automationListening = monitoring === true @@ -1208,7 +1211,33 @@ app.whenReady().then(async () => { ) // 日报等系统动作仍复用现有发送服务;普通聊天不再暴露这个入口。 ipcMain.handle('wechat-personal:send', async (_, request: PersonalWechatSendRequest) => { - if (request.type !== 'voice' || String(request.fromId || '').trim()) { + // 项目规则:**所有发送都必须经过 WechatActionGateway**(审计 + 幂等 + 同一个 Send Log)。 + // 手动发送在改造前从这里直连 `PersonalWechatSendService`,于是完全不留痕 —— + // 排查「消息到底发没发出去」时恰好缺的就是那份证据。 + // `triggerType: 'user'` 不会命中 automation 专用的 purpose allowlist 与 3s 节流, + // 所以行为与改造前一致,只是多一条审计记录(「发送日志」里能看到)。 + if (request.type === 'text' || request.type === 'image') { + const to = String(request.to || '').trim() + const action = await wechatActionGateway.execute({ + origin: 'user_manual', + purpose: request.type === 'text' ? 'manual_text' : 'manual_image', + triggerType: 'user', + recipient: { + type: request.isGroup || to.endsWith('@chatroom') ? 'group' : 'contact', + id: to + }, + content: + request.type === 'text' + ? { type: 'text', text: String(request.text || '') } + : { type: 'image', path: String(request.filePath || '') } + }) + // 返回契约保持 `PersonalWechatSendResult`,界面判读不用改。 + return toPersonalWechatSendResult(action, await personalWechatSendService.getStatus()) + } + + // 语音仍走既有分支:它有自己的网关入口 `wechat-personal:sendGeneratedTtsVoice`, + // 且需要先解析当前账号 wxid 才能编码。 + if (String(request.fromId || '').trim()) { return personalWechatSendService.send(request) } let fromId = '' @@ -2368,11 +2397,13 @@ app.whenReady().then(async () => { if (client) { voiceService = new VoiceService(client, client.getAccountRoot()) voiceRecognition?.connect(voiceService, client.getAccountRoot()) - const monitoring = await client.startMonitor((type, json) => { + const monitoring = await client.startMonitor((event) => { client.invalidateSessionCache() - groupExitMonitorService.notifyDatabaseChanged(json) + groupExitMonitorService.notifyDatabaseChanged(event.raw) for (const window of BrowserWindow.getAllWindows()) { - if (!window.isDestroyed()) window.webContents.send('wcdb-change', { type, json }) + if (!window.isDestroyed()) { + window.webContents.send('wcdb-change', { type: event.type, json: event.raw }) + } } }) void groupExitMonitorService.start(monitoring) diff --git a/src/main/services/automation-action-runner.ts b/src/main/services/automation-action-runner.ts index f01c1a4..524a8a4 100644 --- a/src/main/services/automation-action-runner.ts +++ b/src/main/services/automation-action-runner.ts @@ -4,6 +4,7 @@ import { AUTOMATION_SEND_PURPOSE, DEFAULT_REPLY_TEXT, automationIdempotencyKey, + normalizeReplyDelaySeconds, type AutomationAction, type AutomationRule, type AutomationStep, @@ -31,6 +32,22 @@ import { wechatActionGateway } from './wechat-action-gateway' /** 一步执行完毕后,后续步骤的处置方式。 */ const STEP_ORDER: AutomationStepKey[] = ['received', 'matched', 'reply', 'report', 'send'] +/** + * 前置步骤失败时,后续步骤的 `skipReason`。 + * + * 必须写清「是因为前面那步没成」,否则用户看到一连串「已跳过」会以为是规则没配好。 + */ +const SKIPPED_AFTER_REPLY_FAILURE = '前置步骤失败(回复确认未成功),本次不再继续。' +const SKIPPED_AFTER_REPORT_FAILURE = '前置步骤失败(日报未生成),没有图片可发送。' + +/** + * 规则启用了「发送日报图片」,但手上没有图片文件。 + * + * 这不是「正常跳过」,也不允许退而求其次去发空路径 / 上一张旧图 / 不存在的文件 —— + * 发错东西比不发更糟,所以直接判失败。 + */ +const SEND_WITHOUT_IMAGE_ERROR = '日报图片未生成,无法发送。' + export interface AutomationRunInput { executionId: string rule: AutomationRule @@ -52,6 +69,8 @@ export interface AutomationActionRunnerDependencies { generateReport?: (request: AgentGroupReportRequest) => Promise executeAction?: (request: WechatActionRequest) => Promise now?: () => number + /** 延迟实现。默认真 sleep;单测注入即时 resolve 的假实现,避免真的等 2 秒。 */ + delay?: (ms: number) => Promise } /** 策略层的错误码 → 用户可读短句。UI 直接展示这些文案,不做二次翻译。 */ @@ -86,8 +105,9 @@ function markFailed(step: AutomationStep, at: number, error: string): void { step.error = error } -function markSkipped(step: AutomationStep): void { +function markSkipped(step: AutomationStep, skipReason?: string): void { step.status = 'skipped' + if (skipReason) step.skipReason = skipReason } function actionErrorMessage(code: string | undefined, fallback: string | undefined): string { @@ -102,11 +122,15 @@ export class AutomationActionRunner { private readonly generateReport: (request: AgentGroupReportRequest) => Promise private readonly executeAction: (request: WechatActionRequest) => Promise private readonly now: () => number + private readonly delay: (ms: number) => Promise constructor(dependencies: AutomationActionRunnerDependencies = {}) { this.generateReport = dependencies.generateReport ?? generateAgentGroupReport this.executeAction = dependencies.executeAction ?? ((request) => wechatActionGateway.execute(request)) this.now = dependencies.now ?? (() => Date.now()) + this.delay = + dependencies.delay ?? + ((ms) => new Promise((resolve) => setTimeout(resolve, ms))) } async run(input: AutomationRunInput): Promise { @@ -122,6 +146,22 @@ export class AutomationActionRunner { const sendAction = findAction(input.rule, 'sendReportImage') // ---- 步骤 3:回复确认 ---- + // + // 回复等待:规则一命中就秒回,看起来就是个机器人(消息刚到、回复就到)。 + // 等待时长是**规则自己的一项执行参数**(`rule.replyDelaySeconds`,在 + // 「编辑自动化 → 3 · 触发后执行」里配),所以不同规则可以不一样。 + // + // 等待刻意放在 `reply` 步骤计时**之外** —— `reply.durationMs` 只应该反映发送本身, + // 否则用户看到「回复确认 2000ms」会误以为是发送慢。 + // 已经命中就不再回头重判规则:等待窗口里规则被停用/删掉也不中断本次执行, + // 与 cooldown、同消息幂等的口径一致(都是「命中那一刻」的快照)。 + if (replyAction) { + // 用共享的归一化函数,而不是 `Number(...) || 0`:旧版 rules.json 里没有这个字段, + // 那应该按**默认 2 秒**处理(否则「默认 2 秒」要等用户手动进编辑页才会生效)。 + const replyDelayMs = normalizeReplyDelaySeconds(input.rule.replyDelaySeconds) * 1_000 + if (replyDelayMs > 0) await this.delay(replyDelayMs) + } + if (!replyAction) { markSkipped(stepAt('reply')) } else { @@ -134,7 +174,7 @@ export class AutomationActionRunner { }) if (!sent.ok) { markFailed(replyStep, this.now(), sent.error || '回复确认失败') - markRemainingSkipped(steps, 'reply') + markRemainingSkipped(steps, 'reply', SKIPPED_AFTER_REPLY_FAILURE) return { steps, status: 'failed', errorSummary: replyStep.error } } markSuccess(replyStep, this.now()) @@ -159,7 +199,7 @@ export class AutomationActionRunner { } if (!result.success || !result.pngPath) { markFailed(reportStep, this.now(), result.error || '日报生成失败') - markRemainingSkipped(steps, 'report') + markRemainingSkipped(steps, 'report', SKIPPED_AFTER_REPORT_FAILURE) return { steps, status: 'failed', errorSummary: reportStep.error } } pngPath = result.pngPath @@ -167,12 +207,23 @@ export class AutomationActionRunner { } // ---- 步骤 5:发送日报图片 ---- - // 日报被跳过(或没产出图片)时,这一条也必须 skipped —— 不能凭空发图。 - if (!sendAction || !pngPath) { + // + // 三种情况必须分开判断 —— 合并成 `!sendAction || !pngPath` 会把 + // 「规则要求发图、但图根本没生成」当成正常跳过,execution 还记成 success: + // + // A. 规则**本来就没有启用**这个动作 → skipped,这是正常的,不影响整体结果; + // B. 启用了,但要发的东西不存在 → **failed**,不能假装成功, + // 更不能退而求其次去发空路径 / 上一次的旧图 / 不存在的文件; + // C. 前置(生成日报)已经失败 → 上面就 return 了,走不到这里。 + if (!sendAction) { markSkipped(stepAt('send')) return { steps, status: 'success', ...(pngPath ? { pngPath } : {}) } } const sendStep = stepAt('send') + if (!pngPath) { + markFailed(sendStep, this.now(), SEND_WITHOUT_IMAGE_ERROR) + return { steps, status: 'failed', errorSummary: sendStep.error } + } sendStep.status = 'running' sendStep.startedAt = this.now() const sent = await this.sendThroughGateway(input, 'report', { type: 'image', path: pngPath }) @@ -226,11 +277,15 @@ function findAction(rule: AutomationRule, type: AutomationAction['type']): Autom } /** 把 `after` 之后的步骤全部标成 `skipped`(前一步挂了,后面的不许再动)。 */ -function markRemainingSkipped(steps: AutomationStep[], after: AutomationStepKey): void { +function markRemainingSkipped( + steps: AutomationStep[], + after: AutomationStepKey, + skipReason?: string +): void { const from = STEP_ORDER.indexOf(after) + 1 for (const key of STEP_ORDER.slice(from)) { const step = steps.find((item) => item.key === key) - if (step) markSkipped(step) + if (step) markSkipped(step, skipReason) } } diff --git a/src/main/services/automation-execution-log-service.ts b/src/main/services/automation-execution-log-service.ts index 37c904c..2ee8ae0 100644 --- a/src/main/services/automation-execution-log-service.ts +++ b/src/main/services/automation-execution-log-service.ts @@ -22,6 +22,8 @@ export interface AutomationExecutionLogDependencies { userDataPath?: () => string } +const EXECUTION_STATUSES: AutomationExecution['status'][] = ['running', 'success', 'failed'] + function normalizeExecution(value: unknown): AutomationExecution | null { if (!value || typeof value !== 'object') return null const record = value as Partial @@ -33,7 +35,10 @@ function normalizeExecution(value: unknown): AutomationExecution | null { ruleName: String(record.ruleName || ''), triggerTime: Number(record.triggerTime) || 0, sourceDisplayName: String(record.sourceDisplayName || ''), - status: record.status === 'success' || record.status === 'failed' ? record.status : 'running', + // 不认识的 status(含历史遗留值)一律降级成 `running`,绝不凭空造出成功/失败。 + status: EXECUTION_STATUSES.includes(record.status as AutomationExecution['status']) + ? (record.status as AutomationExecution['status']) + : 'running', durationMs: Number(record.durationMs) || 0, steps: Array.isArray(record.steps) ? record.steps : [], ...(record.errorSummary ? { errorSummary: String(record.errorSummary) } : {}) @@ -78,7 +83,12 @@ export class AutomationExecutionLogService { return true } - /** 统计 `sinceMs` 之后(含)的记录数与成功数。用于顶部「今日执行」。 */ + /** + * 统计 `sinceMs` 之后(含)的记录数与成功数。用于顶部「今日执行」。 + * + * 每条记录都对应一次**真正跑过**的执行(gate 拦下的消息不会产生记录), + * 所以这里直接计数即可。 + */ countSince(sinceMs: number): { total: number; success: number } { this.ensureLoaded() const from = Number(sinceMs) || 0 diff --git a/src/main/services/automation-service.ts b/src/main/services/automation-service.ts index 955a0cf..9f45bb5 100644 --- a/src/main/services/automation-service.ts +++ b/src/main/services/automation-service.ts @@ -52,6 +52,28 @@ import { * 少了它就会自己触发自己,形成死循环。 */ +/** + * gate 表的容量上限。超了才清理「已解锁」的条目 —— 不是为了省内存, + * 而是防止长期运行下 ruleId × conversationId 无界增长。 + */ +const GATE_MAX_ENTRIES = 5_000 + +/** + * 一条规则在一个会话里的 gate 状态。 + * + * ``` + * blocked = inFlight || now < triggeredAt + cooldown + * ``` + */ +interface ConversationRuleGate { + /** **第一次**触发的时间(cooldown 起点)。执行完成时**不重置**。 */ + triggeredAt: number + /** 当前是否有这条规则在这个会话里的执行还在跑。 */ + inFlight: boolean + /** 本次执行期间被 gate 挡下的消息数(仅诊断日志用)。 */ + blockedSinceTrigger: number +} + /** 幂等登记表的存活时间与容量上限(与 MessageListener 的 dedup 同思路)。 */ const CLAIM_TTL_MS = 10 * 60 * 1000 const CLAIM_MAX_ENTRIES = 5_000 @@ -66,6 +88,8 @@ export interface AutomationServiceDependencies { /** 由主进程注入:MessageListener 是否在运行。 */ isListening?: () => boolean now?: () => number + /** gate 表容量上限。默认 `GATE_MAX_ENTRIES`;单测用小值来触发清理路径。 */ + gateMaxEntries?: number } interface ClaimEntry { @@ -82,8 +106,18 @@ export class AutomationService { /** `ruleId:sessionId:localId` → 登记时间。 */ private readonly claims = new Map() - /** `${ruleId}:${conversationId}` → 上次真正开始执行的毫秒时间戳。 */ - private readonly cooldowns = new Map() + /** + * Automation rule ↔ conversation gate。 + * + * 粒度是 **`ruleId + conversationId`**(不是 sender,也不是整个会话): + * 群 A 的「@我生成日报」被触发后,群 B 的同一条规则、或同群里的**另一条规则**都不受影响。 + */ + private readonly gates = new Map() + + /** 仅用于诊断统计(`[Automation] blockedMessages=N`),**不进用户执行日志**。 */ + private blockedMessageCount = 0 + + private readonly gateMaxEntries: number private selfUsernames: string[] = [] private selfUsernamesAt = 0 @@ -98,6 +132,7 @@ export class AutomationService { this.getCapability = dependencies.getCapability ?? getPersonalWechatSendCapability this.isListening = dependencies.isListening ?? (() => true) this.now = dependencies.now ?? (() => Date.now()) + this.gateMaxEntries = dependencies.gateMaxEntries ?? GATE_MAX_ENTRIES } /** @@ -125,6 +160,30 @@ export class AutomationService { } for (const rule of rules) { + // ① Gate 检查放在 TriggerMatcher **之前**(纯读,无副作用)。 + // + // Automation rule conversation gate: + // Once a rule is triggered for a conversation, ignore subsequent messages for + // this rule while the current execution is still running OR until the configured + // cooldown has elapsed from the original trigger time. + // + // Messages arriving during this blocked window must not be matched, replied to, + // executed, or written to the user-facing execution log. + // + // The rule becomes eligible again only after BOTH: + // 1. the previous execution has finished; and + // 2. the cooldown since the original trigger has expired. + // + // Important: cooldown starts at the original trigger time, not when execution finishes. + // + // 也就是说:第一个任务执行期间,以及配置的触发间隔内,后续消息全部静默忽略。 + if (this.isGated(rule, message.sessionId)) { + this.blockedMessageCount += 1 + this.noteBlocked(rule, message.sessionId) + continue + } + + // ② 匹配(同步)。放在 gate 之后 —— 阻塞窗口内的消息不该被匹配。 let matched: boolean try { matched = matchAutomationRule(rule, input, selfUsernames).matched @@ -134,13 +193,28 @@ export class AutomationService { } if (!matched) continue - // 三道闸全部是**同步**的,必须在任何 await 之前完成登记, - // 否则同一批并发消息会同时穿过检查。 - if (this.isCoolingDown(rule, message.sessionId)) continue - if (!this.claim(rule, message)) continue + // ③ Gate claim:check → 写入。**中间不能有 await**,否则两条几乎同时到达的消息 + // 会一起穿过检查。上面 ① 到这里的唯一代码是同步的 matcher,所以这里仍是原子的。 + if (!this.claimGate(rule, message.sessionId)) { + this.blockedMessageCount += 1 + this.noteBlocked(rule, message.sessionId) + continue + } + + // ④ 同一条消息只处理一次(native 事件重复 / 多窗口回读都可能重复投递)。 + if (!this.claim(rule, message)) { + this.leaveGate(rule, message.sessionId) + continue + } const sourceDisplayName = this.resolveDisplayName(message) - await this.execute(rule, message, sourceDisplayName) + try { + await this.execute(rule, message, sourceDisplayName) + } finally { + // 无论 success / failed / 抛异常都要解除 in-flight。 + // 漏掉这一步会让这条规则在这个会话里**永久锁死**。 + this.leaveGate(rule, message.sessionId) + } } } @@ -223,7 +297,8 @@ export class AutomationService { return } - this.cooldowns.set(`${rule.id}:${message.sessionId}`, this.now()) + // cooldown 起点已经在 `claimGate()` 记过了(= 第一次触发时间)。 + // 这里**不再**写任何时间戳 —— 否则「执行真正开始」的时间会把 cooldown 起点往后推。 let result: Awaited> try { @@ -322,15 +397,108 @@ export class AutomationService { } } - private isCoolingDown(rule: AutomationRule, conversationId: string): boolean { + private gateKey(rule: AutomationRule, conversationId: string): string { + // 粒度:规则 × 会话。**不含 senderId** —— 产品语义是「这个群里这条规则刚被触发过」, + // 不是「张三 60 秒内不能再触发」。 + return `${rule.id}:${conversationId}` + } + + private cooldownMs(rule: AutomationRule): number { const seconds = Number(rule.cooldownSeconds) - if (!Number.isFinite(seconds) || seconds <= 0) return false - const last = this.cooldowns.get(`${rule.id}:${conversationId}`) - if (last === undefined) return false - return this.now() - last < seconds * 1000 + if (!Number.isFinite(seconds) || seconds <= 0) return 0 + return seconds * 1000 + } + + /** + * 是否处于阻塞窗口。**纯读**,所以可以放在 TriggerMatcher 之前。 + * + * ``` + * blocked = inFlight || now < triggeredAt + cooldown + * ``` + * + * `triggeredAt` 是**第一次触发**的时间,不是执行完成时间: + * 任务跑得比 cooldown 久时,真正的解锁时间是 + * `max(执行完成, triggeredAt + cooldown)`,而不是「完成之后再等一个 cooldown」。 + */ + private isGated(rule: AutomationRule, conversationId: string): boolean { + const gate = this.gates.get(this.gateKey(rule, conversationId)) + if (!gate) return false + if (gate.inFlight) return true + const cooldownMs = this.cooldownMs(rule) + if (cooldownMs <= 0) return false + // 用「当前时间 vs triggeredAt」现算,**不依赖 setTimeout** —— + // 事件循环卡顿 / app suspend / 定时器漂移都不会把冷却算错。 + return this.now() < gate.triggeredAt + cooldownMs + } + + /** + * 原子地占用 gate(check + 写入)。 + * + * 返回 `false` 表示这一刻已经被阻塞(调用方按「忽略」处理,不匹配、不执行、不写日志)。 + * 调用方必须保证**从 `isGated()` 到这里之间没有 await**。 + */ + private claimGate(rule: AutomationRule, conversationId: string): boolean { + if (this.isGated(rule, conversationId)) return false + this.gates.set(this.gateKey(rule, conversationId), { + triggeredAt: this.now(), + inFlight: true, + blockedSinceTrigger: 0 + }) + this.evictGatesIfNeeded() + return true + } + + /** + * 解除 in-flight。**保留 `triggeredAt`** —— cooldown 仍以第一次触发时间为起点。 + * + * 无论 success / failed / 抛异常都必须走这里,否则这条规则在这个会话里会永久锁死。 + */ + private leaveGate(rule: AutomationRule, conversationId: string): void { + const gate = this.gates.get(this.gateKey(rule, conversationId)) + if (!gate) return + gate.inFlight = false + // 本次执行期间被挡下的消息数,汇总成**一行**诊断日志。 + // 只允许出现数量与 ruleId(不含群名 / wxid / 昵称 / 正文)。 + if (gate.blockedSinceTrigger > 0) { + this.info(`blockedMessages=${gate.blockedSinceTrigger} ruleId=${rule.id}`) + gate.blockedSinceTrigger = 0 + } + } + + /** 记一次「被 gate 拦下的消息」。**只计数**,绝不写入用户执行日志。 */ + private noteBlocked(rule: AutomationRule, conversationId: string): void { + const gate = this.gates.get(this.gateKey(rule, conversationId)) + if (gate) gate.blockedSinceTrigger += 1 + } + + /** 仅供诊断:累计被 gate 拦下的消息数。不接 IPC、不进 UI。 */ + getBlockedMessageCount(): number { + return this.blockedMessageCount + } + + /** + * gate 表有界:**只清理已经解锁的条目**。 + * + * 硬不变量:**正在执行(`inFlight === true`)的 gate 永不被淘汰** —— + * 淘汰它就等于把一条正在跑的规则提前放开,会立刻产生第二次并发执行。 + * 所以这里的条件是 `!inFlight && 已过 cooldown`,两个都要满足; + * 清理不掉就一直留着(宁可让表大一点,也不能让正在跑的规则失去保护)。 + */ + private evictGatesIfNeeded(): void { + if (this.gates.size <= this.gateMaxEntries) return + const now = this.now() + const rules = this.ruleStore.listRules() + for (const [key, gate] of this.gates) { + if (this.gates.size <= this.gateMaxEntries) break + if (gate.inFlight) continue + const separator = key.lastIndexOf(':') + const ruleId = separator < 0 ? key : key.slice(0, separator) + const rule = rules.find((item) => item.id === ruleId) + const cooldownMs = rule ? this.cooldownMs(rule) : 0 + if (now >= gate.triggeredAt + cooldownMs) this.gates.delete(key) + } } - /** 登记「这条消息已经被这条规则处理过」。返回 false 表示重复,应当跳过。 */ private claim(rule: AutomationRule, message: NormalizedIncomingMessage): boolean { const key = `${rule.id}:${message.sessionId}:${message.localId}` const now = this.now() diff --git a/src/main/services/message-listener-service.ts b/src/main/services/message-listener-service.ts index 7eeaa37..007b196 100644 --- a/src/main/services/message-listener-service.ts +++ b/src/main/services/message-listener-service.ts @@ -1,5 +1,5 @@ import { decompress as zstdDecompress } from 'fzstd' -import type { Wcdb4Client, Wcdb4Message } from '../wcdb4-client' +import type { Wcdb4Client, Wcdb4Message, Wcdb4MonitorEvent } from '../wcdb4-client' /** * MessageListener —— 实时消息回读底座。 @@ -21,9 +21,8 @@ import type { Wcdb4Client, Wcdb4Message } from '../wcdb4-client' * **不负责**(这些属于下一层,本模块不许碰):关键词匹配、@我业务判断、日报、 * 自动回复、AI 调用、Agent 调度、发送消息。 * - * 设计约束全部来自 2026-09-20 的真实环境 Spike(见 - * `.ai-local/reports/2026-09-20-realtime-message-spike-result.md`): - * 一条真实消息平均触发 **≈17 个** native event,所以**绝不能**一个事件触发一次业务动作。 + * 设计约束:一次写入会触发**一连串** native event(十几条), + * 所以**绝不能**一个事件触发一次业务动作。 */ /** zstd 帧魔数(微信 `source` 列是 zstd 压缩)。 */ @@ -185,6 +184,13 @@ export class MessageListenerService { * 必须有界 —— Spike 里用的是无上限 `Set`,长时间运行会持续吃内存。 */ private readonly seen = new Map() + /** + * 本批 coalesce 窗口内**出现过精确会话**的事件收集到的 session 集合。 + * + * v2 事件(Native Monitor Event v2)会带上发生变化的 sessionId,这里用 **Set** 收集 —— + * 120ms 内收到 `A B A C B` 必须回读 A/B/C 三个,绝不能「最后一个 wins」。 + */ + private readonly pendingSessions = new Set() private coalesceTimer: ReturnType | null = null private readbackInFlight = false private disposed = false @@ -243,11 +249,17 @@ export class MessageListenerService { /** * 接 native change event。 * - * 只做计数与 coalesce —— **绝不**在这里回读,更不触发任何业务。 + * 只做计数、收集会话与 coalesce —— **绝不**在这里回读,更不触发任何业务。 + * + * `event.protocol === 2` 时把 `sessionId` 收进本批的集合;legacy 事件(v1,不带会话) + * 走到回读时仍然只能退化成「最近活跃会话」。 */ - handleNativeChange(): void { + handleNativeChange(event?: Wcdb4MonitorEvent): void { if (this.disposed) return this.counters.nativeEvents += 1 + if (event?.protocol === 2 && event.sessionId) { + this.pendingSessions.add(event.sessionId) + } if (this.coalesceTimer) { // 已经在同一个窗口里:直接合并掉,这就是 17:1 那个比值的收敛点。 this.counters.coalescedEvents += 1 @@ -262,47 +274,44 @@ export class MessageListenerService { /** * 有界回读 → normalize → dedup → 投递。 * - * 回读范围严格受限:**单会话** + **小时间窗口** + **条数上限**。 - * 不遍历会话、不扫全库。 + * 回读范围严格受限:**按会话** + **小时间窗口** + **条数上限**。不遍历会话、不扫全库。 + * + * 目标会话的来源有两种: + * - **precise(protocol v2)**:本批事件收集到的 `sessionId` 集合,逐个回读 + * (多个会话同时来消息时,A/B/C 都会被读到); + * - **legacy(protocol v1,旧 runtime)**:只能退回「最近活跃会话」`getSessions()[0]`。 */ private async readback(): Promise { if (this.disposed || this.readbackInFlight) return this.readbackInFlight = true const startedAt = Date.now() + + const preciseSessions = Array.from(this.pendingSessions) + this.pendingSessions.clear() + try { - const nowSec = Math.floor(startedAt / 1000) - const sessions = this.client.getSessions() - const session = sessions[0] - if (!session?.username) return - - this.counters.readbacks += 1 - const messages = await this.client.getMessagesAsync( - session.username, - nowSec - this.lookbackSec, - nowSec + this.lookaheadSec, - { limit: this.readLimit } - ) - if (this.disposed) return - let delivered = 0 - for (const message of messages) { - const normalized = this.normalize(session.username, message) - if (!normalized) continue - // dedup:同一条消息可能在多个读回窗口中反复出现。 - if (!this.markSeen(normalized)) { - this.counters.deduped += 1 - continue + let sessions = 0 + + if (preciseSessions.length > 0) { + for (const sessionId of preciseSessions) { + if (this.disposed) return + delivered += await this.readbackSession(sessionId) + sessions += 1 } - this.counters.delivered += 1 - delivered += 1 - this.deliver(normalized) + } else { + // legacy:旧 runtime 的 payload 不带会话,只能回读「最近活跃会话」。 + const session = this.client.getSessions()[0] + if (!session?.username) return + delivered = await this.readbackSession(session.username) + sessions = 1 } if (delivered > 0) { - // 日志只记录数量与耗时,**不含** wxid / 群名 / 昵称 / 正文 / source。 + // 日志只记录协议版本、数量与耗时 —— **不含** wxid / 群名 / 昵称 / 正文 / source / sessionId。 console.log( - `[MessageListener] delivered=${delivered} readbackMs=${Date.now() - startedAt}` + - ` sessionType=${session.username.endsWith('@chatroom') ? 'group' : 'direct'}` + + `[MessageListener] protocol=${preciseSessions.length > 0 ? 'v2' : 'v1'}` + + ` sessions=${sessions} delivered=${delivered} readbackMs=${Date.now() - startedAt}` + ` events=${this.counters.nativeEvents} coalesced=${this.counters.coalescedEvents}` ) } @@ -315,6 +324,34 @@ export class MessageListenerService { } } + /** 回读**单个**会话并投递。返回本次真正投递出去的条数。 */ + private async readbackSession(sessionId: string): Promise { + const nowSec = Math.floor(Date.now() / 1000) + this.counters.readbacks += 1 + const messages = await this.client.getMessagesAsync( + sessionId, + nowSec - this.lookbackSec, + nowSec + this.lookaheadSec, + { limit: this.readLimit } + ) + if (this.disposed) return 0 + + let delivered = 0 + for (const message of messages) { + const normalized = this.normalize(sessionId, message) + if (!normalized) continue + // dedup:同一条消息可能在多个读回窗口中反复出现(native 事件本身也可能重复)。 + if (!this.markSeen(normalized)) { + this.counters.deduped += 1 + continue + } + this.counters.delivered += 1 + delivered += 1 + this.deliver(normalized) + } + return delivered + } + private deliver(message: NormalizedIncomingMessage): void { for (const listener of this.listeners) { try { diff --git a/src/main/services/personal-wechat-send-service.ts b/src/main/services/personal-wechat-send-service.ts index 3f03874..3ad4ecf 100644 --- a/src/main/services/personal-wechat-send-service.ts +++ b/src/main/services/personal-wechat-send-service.ts @@ -931,6 +931,44 @@ function windowsStatusBase(host = windowsHookHost() || ''): PersonalWechatSender } } +/** 文件字节数;读不到时返回 `-1`(**不是** 0 —— 0 会被误读成「空文件」)。 */ +function safeFileSize(target: string): number { + try { + return statSync(target).size + } catch { + return -1 + } +} + +/** + * 发送链路诊断日志(排查「hook 报成功、群里却没有消息」时用)。 + * + * 只有在设置里打开「调试模式」时才输出,所以正常使用不会刷屏。 + * + * **隐私约束**:只允许输出类型、字节数、HTTP 状态码、hook 的 `ret` 这类结构性字段。 + * 文件名、完整路径、wxid、群名、消息正文不进日志。 + */ +function logSendDiagnostic(stage: string, detail: Record): void { + try { + if (!loadSettings().debugEnabled) return + } catch { + return + } + console.log(`[SendDiag] ${stage} ${JSON.stringify(detail)}`) +} + +/** + * 发送完成后**保留**图片临时文件的时长(毫秒)。默认 `0` = 立刻删除(现行行为)。 + * + * 为什么留这个开关:`finally` 里的删除发生在 **hook 返回之后**,这隐含假设 + * 「hook 返回时发送已经完成」。如果 hook 其实是「先回 200、后台再读文件」, + * 文件就可能在它读到之前已经被删掉 —— 现象正好是「hook 报成功、微信里没有消息」。 + * + * 置为非 0 时可分辨两类失败:保留期间图片到了 ⇒ 删除时机问题(修法与 hook 约定文件生命周期); + * 仍然没到 ⇒ 与文件生命周期无关。 + */ +const WINDOWS_TEMP_IMAGE_GRACE_MS = 0 + async function requestWindowsHook( endpoint: string, body: Record, @@ -939,15 +977,41 @@ async function requestWindowsHook( host = windowsHookHost() ): Promise { if (!host) throw new Error('尚未配置微信发送能力端口') + const startedAt = Date.now() const init: RequestInit = { method, headers: { 'Content-Type': 'application/json' } } if (method !== 'GET') init.body = JSON.stringify(body) + logSendDiagnostic('hook-request', { + endpoint, + method, + // 只报 payload 里有没有内容,不报内容本身。 + hasPayload: method !== 'GET' ? init.body !== undefined : false + }) const response = await requestWithTimeout(`http://${host}${endpoint}`, init, timeoutMs) const responseText = await response.text() - if (!response.ok) throw new WindowsHookHttpError(response.status, responseText) - return parseWindowsHookResponse(responseText, method === 'POST') + if (!response.ok) { + // HTTP 层失败:状态码不是隐私,正文可能是(hook 有时会把请求原样回显)。 + logSendDiagnostic('hook-response', { + endpoint, + httpStatus: response.status, + ok: false, + bodyLength: responseText.length, + elapsedMs: Date.now() - startedAt + }) + throw new WindowsHookHttpError(response.status, responseText) + } + const parsed = parseWindowsHookResponse(responseText, method === 'POST') + logSendDiagnostic('hook-response', { + endpoint, + httpStatus: response.status, + ok: true, + ret: typeof parsed.ret === 'number' ? parsed.ret : null, + retMessageLength: String(parsed.retmsg ?? parsed.msg ?? '').length, + elapsedMs: Date.now() - startedAt + }) + return parsed } export class PersonalWechatSendService { @@ -1465,6 +1529,14 @@ export class PersonalWechatSendService { } else if (request.type === 'image') { const preparedImage = prepareWindowsImageFile(request.filePath) temporaryImagePath = preparedImage.temporary ? preparedImage.filePath : undefined + // 关键证据:源文件有多大、是否走了临时副本、副本拷出来多大。 + // `copiedBytes < sourceBytes` ⇒ 截断的拷贝:拷贝中途失败不会抛错, + // 只看「有没有文件」会漏掉这种情况。 + logSendDiagnostic('image-prepared', { + sourceBytes: safeFileSize(request.filePath), + usedTempCopy: preparedImage.temporary, + copiedBytes: safeFileSize(preparedImage.filePath) + }) const windowsRequest = buildWindowsWechatRequest({ ...request, filePath: preparedImage.filePath @@ -1503,10 +1575,22 @@ export class PersonalWechatSendService { } } finally { if (temporaryImagePath) { - try { - unlinkSync(temporaryImagePath) - } catch { - // Temporary files are best-effort cleanup only. + if (WINDOWS_TEMP_IMAGE_GRACE_MS > 0) { + // 诊断模式:先不删,等一段时间再删(见常量注释)。 + const pending = temporaryImagePath + setTimeout(() => { + try { + unlinkSync(pending) + } catch { + // best effort + } + }, WINDOWS_TEMP_IMAGE_GRACE_MS).unref?.() + } else { + try { + unlinkSync(temporaryImagePath) + } catch { + // Temporary files are best-effort cleanup only. + } } } if (temporaryVoicePath) { diff --git a/src/main/services/wechat-action-gateway.ts b/src/main/services/wechat-action-gateway.ts index bff8161..adef5e5 100644 --- a/src/main/services/wechat-action-gateway.ts +++ b/src/main/services/wechat-action-gateway.ts @@ -5,7 +5,8 @@ import path from 'path' import type { PersonalWechatSendCapability, PersonalWechatSendRequest, - PersonalWechatSendResult + PersonalWechatSendResult, + PersonalWechatSenderStatus } from '../../shared/personal-wechat' import type { PolicyDecision, @@ -608,6 +609,32 @@ function toPersonalWechatSendRequest(request: WechatActionRequest): PersonalWech } } +/** + * 把网关结果还原成既有调用方期望的 `PersonalWechatSendResult`。 + * + * 为什么需要:手动发送的 IPC(`wechat-personal:send`)返回契约是 + * `PersonalWechatSendResult`,界面上靠 `response.success` / `response.error` 判读。 + * 改成走网关之后必须把这个契约**原样**还回去,否则「统一发送路径」会把 UI 一起改坏。 + * + * 网关成功时会把底层 `sendResult` 原样挂在 `action.sendResult` 上,优先用它 + * (它带着真实的 `status`);拿不到时按 `action.status` 合成一个。 + */ +export function toPersonalWechatSendResult( + action: WechatActionResult, + fallbackStatus: PersonalWechatSenderStatus +): PersonalWechatSendResult { + const raw = action.sendResult + if (raw && typeof raw === 'object' && 'success' in (raw as Record)) { + return raw as PersonalWechatSendResult + } + if (action.status === 'sent') return { success: true, status: fallbackStatus } + return { + success: false, + status: fallbackStatus, + error: action.reason || action.errorCode || '微信发送失败' + } +} + function contentAudit( content: WechatActionContent ): Pick { diff --git a/src/main/wcdb4-client.ts b/src/main/wcdb4-client.ts index d111151..4a1c848 100644 --- a/src/main/wcdb4-client.ts +++ b/src/main/wcdb4-client.ts @@ -384,6 +384,79 @@ export function resolveWindowsNativeAccountRoot( } } +/** + * native monitor 事件(pipe / socket 上的一行 JSON)。 + * + * ## 两个协议版本 + * + * **v1(历史)** —— 变更信号来自**文件系统通知**(Windows `ReadDirectoryChangesW` / + * macOS `kqueue`),native 只知道「哪个文件被写了」: + * + * ```json + * {"db":"message_0.db","table":"message","action":"update"} + * ``` + * + * 它**不携带会话**,所以上层只能回读「最近活跃会话」。`protocol` 会标成 1。 + * + * **v2(Native Monitor Event v2)** —— native 侧在收到文件变化后做一次 + * SessionTable watermark diff,把**每一个真正变化的会话**各发一条: + * + * ```json + * {"version":2,"kind":"message_change","session_id":"xxx@chatroom", + * "db":"message_0.db","table":"message","action":"update","observed_at_ms":1} + * ``` + * + * `protocol = 2` 且带 `sessionId` 时,上层可以**直接按会话回读**,不再依赖 + * `getSessions()[0]`。 + * + * **降级规则**:只要 `version`/`kind`/`session_id` 任一不符合 v2 契约, + * 就按 v1 处理(`protocol = 1`)—— 宁可不精确,也不能拿一个不可信的会话 id 去读。 + */ +export interface Wcdb4MonitorEvent { + /** pipe 上的原始一行(原样转发给需要它的消费者,例如退群监控)。 */ + raw: string + /** 语义化的动作类型。当事件语义出现时用它,比如 `'user_change'`。 */ + type: string + action: string + /** `2` = 精确会话事件;`1` = legacy(不含会话)。 */ + protocol: 1 | 2 + /** 粗表分类:`database` / `message` / `Session` / `contact`。 */ + table: string + /** `protocol === 2` 时才有:发生变化的会话(`xxx@chatroom` 或 wxid)。 */ + sessionId?: string + /** native 侧观测时刻(epoch ms),仅用于诊断。 */ + observedAtMs?: number +} + +/** 把 pipe 上的一行 payload 解析成事件。纯函数,便于单测。 */ +export function parseMonitorEvent(rawPayload: string): Wcdb4MonitorEvent { + const raw = rawPayload.trim() + let parsed: Record = {} + try { + const value = JSON.parse(raw) as unknown + if (value && typeof value === 'object' && !Array.isArray(value)) { + parsed = value as Record + } + } catch { + // 不是 JSON:当作 legacy 事件处理,保留原始 payload。 + } + + const action = typeof parsed.action === 'string' && parsed.action ? parsed.action : 'update' + const sessionId = typeof parsed.session_id === 'string' ? parsed.session_id.trim() : '' + // 三个条件缺一不可:版本、语义、以及**非空**的会话 id。 + const precise = parsed.version === 2 && parsed.kind === 'message_change' && Boolean(sessionId) + + return { + raw, + type: action, + action, + protocol: precise ? 2 : 1, + table: typeof parsed.table === 'string' ? parsed.table : '', + ...(precise ? { sessionId } : {}), + ...(typeof parsed.observed_at_ms === 'number' ? { observedAtMs: parsed.observed_at_ms } : {}) + } +} + export class Wcdb4Client { static readonly defaultRoot = Wcdb4Client.findExistingDefaultRoot() @@ -495,7 +568,7 @@ export class Wcdb4Client { private wcdbStopMonitorPipe: (() => void) | null = null private wcdbGetMonitorPipeName: ((outName: WcdbVoidOut) => number) | null = null private monitorPipeClient: Socket | null = null - private monitorCallback: ((type: string, json: string) => void) | null = null + private monitorCallback: ((event: Wcdb4MonitorEvent) => void) | null = null private monitorConnectTimer: ReturnType | null = null private monitorReconnectTimer: ReturnType | null = null private monitorPipePath = '' @@ -773,7 +846,7 @@ export class Wcdb4Client { }) } - async startMonitor(callback: (type: string, json: string) => void): Promise { + async startMonitor(callback: (event: Wcdb4MonitorEvent) => void): Promise { if (this.closing || !this.wcdbStartMonitorPipe || !this.wcdbGetMonitorPipeName || !this.koffi) { return false } @@ -911,13 +984,7 @@ export class Wcdb4Client { private emitMonitorPayload(rawPayload: string): void { const payload = rawPayload.trim() if (!payload || !this.monitorCallback) return - - try { - const parsed = JSON.parse(payload) as { action?: string } - this.monitorCallback(parsed.action || 'update', payload) - } catch { - this.monitorCallback('update', payload) - } + this.monitorCallback(parseMonitorEvent(payload)) } private scheduleMonitorReconnect(): void { diff --git a/src/renderer/src/features/automation/AutomationWorkspace.tsx b/src/renderer/src/features/automation/AutomationWorkspace.tsx index 4081fc0..d94b2f3 100644 --- a/src/renderer/src/features/automation/AutomationWorkspace.tsx +++ b/src/renderer/src/features/automation/AutomationWorkspace.tsx @@ -9,9 +9,6 @@ import { Button, SegmentedControl, SegmentedControlItem, - Tooltip, - TooltipContent, - TooltipTrigger, useToast } from '../../components/ui' import { RuleListPanel } from './RuleListPanel' @@ -188,21 +185,8 @@ export function AutomationWorkspace({
消息监听 - {status.listening ? '正在监听新消息' : '未在监听'} + {status.listening ? '运行中' : '未在监听'} - {status.listeningDegraded ? ( - - - - 当前仅能捕获最近活跃的会话 - - - - 微信底层只会上报「有表发生了变化」,不带是哪条会话。TraceMemo - 目前据此回读最近活跃的会话,因此同一时刻多个群同时来消息时,只会处理其中的一个。 - - - ) : null}
diff --git a/src/renderer/src/features/automation/ExecutionDetailDrawer.tsx b/src/renderer/src/features/automation/ExecutionDetailDrawer.tsx index 2585ef7..139ed97 100644 --- a/src/renderer/src/features/automation/ExecutionDetailDrawer.tsx +++ b/src/renderer/src/features/automation/ExecutionDetailDrawer.tsx @@ -1,4 +1,5 @@ import * as React from 'react' +import { AUTOMATION_EXECUTION_STATUS_LABELS } from '../../../../shared/automation' import type { AutomationExecution } from '../../../../shared/automation' import { Button } from '../../components/ui' import { @@ -75,7 +76,7 @@ export function ExecutionDetailDrawer({
结果
- {execution.status === 'success' ? '成功' : '失败'} + {AUTOMATION_EXECUTION_STATUS_LABELS[execution.status]}
diff --git a/src/renderer/src/features/automation/ExecutionLogPanel.tsx b/src/renderer/src/features/automation/ExecutionLogPanel.tsx index 2cd5bf5..a5267d4 100644 --- a/src/renderer/src/features/automation/ExecutionLogPanel.tsx +++ b/src/renderer/src/features/automation/ExecutionLogPanel.tsx @@ -1,4 +1,5 @@ import * as React from 'react' +import { AUTOMATION_EXECUTION_STATUS_LABELS } from '../../../../shared/automation' import type { AutomationExecution } from '../../../../shared/automation' import { Button, Spinner } from '../../components/ui' import { STEP_STATUS_LABELS, formatDuration, formatTriggerTime, stepStatusTone } from './model/format' @@ -100,7 +101,7 @@ export function ExecutionLogPanel({ - {execution.status === 'success' ? '成功' : '失败'} + {AUTOMATION_EXECUTION_STATUS_LABELS[execution.status]} diff --git a/src/renderer/src/features/automation/RuleEditorPanel.tsx b/src/renderer/src/features/automation/RuleEditorPanel.tsx index 5e66126..ff0722b 100644 --- a/src/renderer/src/features/automation/RuleEditorPanel.tsx +++ b/src/renderer/src/features/automation/RuleEditorPanel.tsx @@ -1,8 +1,10 @@ import * as React from 'react' import { + AUTOMATION_REPLY_DELAY_MAX_SECONDS, DEFAULT_REPLY_TEXT, KEYWORD_MATCH_MODE_LABELS, createDefaultDailyReportRule, + normalizeReplyDelaySeconds, normalizeRuleDraft, type AutomationActionType, type AutomationRule, @@ -54,7 +56,10 @@ function draftFromRule(rule: AutomationRule | null): AutomationRuleDraft { scope: base.scope, conditions: base.conditions, actions: base.actions, - cooldownSeconds: base.cooldownSeconds + cooldownSeconds: base.cooldownSeconds, + // 漏掉一个字段就会被 normalizeRuleDraft 的默认值悄悄覆盖 —— + // 比如把用户设成 0 的「回复前等待」重置回 2 秒。 + replyDelaySeconds: base.replyDelaySeconds }) } @@ -122,6 +127,15 @@ export function RuleEditorPanel({ ) }, [groups, groupFilter]) + /** + * 生效范围的实际状态。 + * + * 现在的交互是「勾选群 = 指定群聊;一个都不勾 = 所有群聊」,所以没有单独的 + * 「所有群聊 / 指定群聊」单选控件 —— 不去动第 2 节的结构,只把这个状态显式说出来, + * 好让「所有群聊」那条例外提示挂在对的上下文里。 + */ + const hasSelectedGroups = draft.conditions.conversationIds.length > 0 + const handleSave = (): void => { if (!draft.name.trim()) { setValidationError('请填写自动化名称') @@ -228,9 +242,9 @@ export function RuleEditorPanel({

2 · 在哪些聊天生效

- {draft.conditions.conversationIds.length + {hasSelectedGroups ? `已选 ${draft.conditions.conversationIds.length} 个群` - : '未选择时对所有群聊生效'} + : '所有群聊'}
+ {hasSelectedGroups ? null : ( + // 只在「所有群聊」这个上下文里说明真实边界。 + // 不解释底层原因(表事件 / 会话回读),也不写「100% 覆盖」这种绝对承诺。 +

+ + 当前版本在多个会话同时收到消息时,极少数自动化触发可能遗漏。 +

+ )}
@@ -268,6 +290,30 @@ export function RuleEditorPanel({

3 · 触发后执行

按下列顺序执行 + {/* 第一步不是动作,而是「等多久」—— 所以放在动作列表之前,不混进 ACTION_META。 */} +
+
+ 回复前等待(秒) + 命中后先等这么久再回复,避免秒回显得像机器人;0 = 立刻回复 +
+ { + const next = Number(event.target.value) + setDraft((current) => ({ + ...current, + replyDelaySeconds: normalizeReplyDelaySeconds( + Number.isFinite(next) ? next : 0 + ) + })) + }} + className="automation-number-input" + aria-label="回复前等待秒数" + /> +
{ACTION_META.map((meta) => { const action = actionOf(meta.type) return ( @@ -304,7 +350,10 @@ export function RuleEditorPanel({
触发间隔(秒) - 同一个群里,两次触发之间至少间隔这么久,避免刷屏 + {/* 这句就是本产品的门语义定义,不要写 debounce / mutex / in-flight 这类技术词。 */} + + 触发后,在当前任务执行期间及设定间隔内,不再处理本规则在该会话中的新消息。 +
item.key === step.key) const blockedByFailure = steps .slice(0, index < 0 ? 0 : index) @@ -60,6 +64,6 @@ export function describeSkippedStep(step: AutomationStep, steps: AutomationStep[ if (blockedByFailure) return '上一步失败,已跳过' if (step.key === 'reply') return '规则未启用「回复确认」' if (step.key === 'report') return '规则未启用「生成日报」' - if (step.key === 'send') return '没有可发送的日报图片' + if (step.key === 'send') return '规则未启用「发送日报图片」' return '未执行' } diff --git a/src/renderer/src/styles/automation.scss b/src/renderer/src/styles/automation.scss index 6cfc450..b8f4336 100644 --- a/src/renderer/src/styles/automation.scss +++ b/src/renderer/src/styles/automation.scss @@ -147,6 +147,21 @@ font-size: 12px; } +// 轻量边界说明(例如「所有群聊」时的遗漏提示)。 +// 只陈述用户能理解的结果,不解释底层机制,也不做「100% 覆盖」这类承诺。 +.automation-section-hint { + display: flex; + gap: 6px; + margin: 10px 0 0; + padding: 8px 10px; + border-radius: var(--wxex-radius-sm); + border: 1px solid var(--wxex-border); + background: color-mix(in srgb, var(--wxex-text-muted) 6%, transparent); + color: var(--wxex-text-secondary); + font-size: 12px; + line-height: 18px; +} + .automation-loading { display: flex; align-items: center; @@ -855,11 +870,16 @@ color: var(--wxex-success); } - &.failed, - &.running { + &.failed { background: color-mix(in srgb, var(--wxex-danger) 16%, transparent); color: var(--wxex-danger); } + + // 「执行中」用强调色,不是红色。 + &.running { + background: var(--wxex-brand-soft); + color: var(--wxex-brand); + } } .automation-trail { diff --git a/src/shared/automation.ts b/src/shared/automation.ts index dbcf5cc..3dafce8 100644 --- a/src/shared/automation.ts +++ b/src/shared/automation.ts @@ -50,6 +50,17 @@ export interface AutomationRule { actions: AutomationAction[] /** 同一规则在同一会话内的最小触发间隔(秒)。防刷屏。 */ cooldownSeconds: number + /** + * 命中后**等多久再发「回复确认」**(秒)。0 = 立刻回复。 + * + * 与 `cooldownSeconds` 同口径(秒):界面上是秒、落盘也是秒, + * 只有 `AutomationActionRunner` 在等待前换算成毫秒。别在这里存毫秒, + * 否则「触发间隔」和「回复等待」两个字段一个秒一个毫秒,迟早有人读错。 + * + * **旧版 rules.json 里没有这个字段**:读取方一律按 + * `normalizeReplyDelaySeconds()` 处理 ⇒ 缺省即默认 2 秒,不是「不等待」。 + */ + replyDelaySeconds: number createdAt: number updatedAt: number } @@ -57,6 +68,7 @@ export interface AutomationRule { /** 执行步骤状态。`skipped` 用于「上一步失败所以这一步不做」。 */ export type AutomationStepStatus = 'pending' | 'running' | 'success' | 'failed' | 'skipped' +/** 执行步骤键。只有**真正进入执行**的规则才会有步骤。 */ export type AutomationStepKey = 'received' | 'matched' | 'reply' | 'report' | 'send' export interface AutomationStep { @@ -69,8 +81,37 @@ export interface AutomationStep { durationMs?: number /** 用户可读的错误信息。**禁止**含 wxid / 文件路径 / 原始 payload。 */ error?: string + /** + * 被跳过时的原因(给用户看的一句话)。 + * + * 与 `error` 分开:`skipped` 不是错误,硬塞进 `error` 会让界面把它当失败渲染。 + * **禁止**含 localId / serverId / dedup key 这类工程术语。 + */ + skipReason?: string } +/** + * 一次 execution 的最终状态。 + * + * **规则(不要改口径)** + * + * | 情况 | status | + * | --- | --- | + * | 所有**启用的必需 action** 成功 | `success` | + * | 某个必需 action `failed` | `failed` | + * | 后续 action 因前置失败未执行(step 记 `skipped`) | `failed`(整个规则没跑完整) | + * | 规则本来就没启用某个 action | 该 step `skipped`,**不影响**整体结果 | + * | 消息**没有命中**规则 | **不创建 execution** | + * | 落在 rule conversation gate 的阻塞窗口内 | **不创建 execution** | + * + * ⚠️ **阻塞窗口内的消息不留任何痕迹**:不匹配、不执行、不回复、**不写 execution**。 + * 从用户产品视角,那段时间这条规则就是「没处理这些消息」—— 不是「跳过」。 + * 门语义见 `automation-service.ts` 的 `ConversationRuleGate`。 + * + * ⇒ 所以 execution 级**没有** `skipped`:一次 execution 只要被创建出来,就说明它真的跑过, + * 结果非成功即失败。**步骤级**的 `skipped`(`AutomationStepStatus`)是另一回事,仍然保留 —— + * 那是「这次执行里某一步没做(没配 / 前置失败)」。 + */ export type AutomationExecutionStatus = 'running' | 'success' | 'failed' export interface AutomationExecution { @@ -108,6 +149,39 @@ export const AUTOMATION_SEND_PURPOSE = { /** 审计日志里的来源标识(`WechatActionOrigin`)。 */ export const AUTOMATION_SEND_ORIGIN = 'automation' +/** + * 「回复确认」前的默认等待(秒)。 + * + * 规则一命中就立刻回复,看起来就是个机器人 —— 消息刚到、回复就到了。 + * 这个等待把回复推到一个更像人的时间点上,所以它是**规则自己的一项执行参数** + * (和 `cooldownSeconds` 同口径:秒),在「编辑自动化 → 3 · 触发后执行」里配。 + * + * **只作用于「回复确认」这一步**:日报生成 + 图片发送本来就以秒计, + * 不再叠加等待,否则整条链路会慢得让人以为卡住了。 + */ +export const AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS = 2 + +/** 等待上限。再长用户会以为功能坏了,所以封顶而不是无限放开。 */ +export const AUTOMATION_REPLY_DELAY_MAX_SECONDS = 60 + +/** + * 把任意输入收敛成合法的等待秒数。 + * + * 口径(三条都是刻意的): + * - **空值 / 非法值 → 默认值**,而不是 0 —— 填错了不该静默变成「不等待」; + * - **0 是合法值**(显式表示立刻回复),负数与它同归 0; + * - 超过上限按上限截断,避免有人填 10 分钟。 + */ +export function normalizeReplyDelaySeconds(value: unknown): number { + if (value === null || value === undefined || value === '') { + return AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS + } + const parsed = typeof value === 'number' ? value : Number(String(value).trim()) + if (!Number.isFinite(parsed)) return AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS + if (parsed <= 0) return 0 + return Math.min(AUTOMATION_REPLY_DELAY_MAX_SECONDS, Math.round(parsed)) +} + /** 所有自动化行为的 `sourceId` / `executionId` 口径:一次执行一个 id。 */ export function automationIdempotencyKey( kind: keyof typeof AUTOMATION_SEND_PURPOSE, @@ -125,6 +199,8 @@ export interface AutomationRuleDraft { conditions: AutomationConditions actions: AutomationAction[] cooldownSeconds: number + /** 命中后等多久再回复确认(秒)。0 = 立刻回复。 */ + replyDelaySeconds: number } /** 顶部状态条数据(全部来自 main 侧真实能力,UI 不做平台判断)。 */ @@ -197,6 +273,18 @@ export const AUTOMATION_STEP_LABELS: Record = { send: '发送日报图片' } +/** + * 执行结果的展示文案。 + * + * `skipped` 是**中性**状态(命中了但被规则自身的冷却 / 去重拦下), + * 不是失败 —— 界面不能用红色渲染它。 + */ +export const AUTOMATION_EXECUTION_STATUS_LABELS: Record = { + running: '执行中', + success: '成功', + failed: '失败' +} + export const KEYWORD_MATCH_MODE_LABELS: Record = { contains: '包含', exact: '完全匹配', @@ -281,6 +369,7 @@ export function createDefaultDailyReportRule(now: number): AutomationRule { { type: 'sendReportImage', enabled: true } ], cooldownSeconds: 60, + replyDelaySeconds: AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS, createdAt: now, updatedAt: now } @@ -323,6 +412,7 @@ export function normalizeRuleDraft(input: unknown, fallbackName = '未命名自 ignoreSelf: conditions.ignoreSelf !== false }, actions: normalizedActions, - cooldownSeconds: Number.isFinite(cooldown) ? Math.max(0, Math.floor(cooldown)) : 60 + cooldownSeconds: Number.isFinite(cooldown) ? Math.max(0, Math.floor(cooldown)) : 60, + replyDelaySeconds: normalizeReplyDelaySeconds(raw.replyDelaySeconds) } } diff --git a/src/shared/wechat-action.ts b/src/shared/wechat-action.ts index 423ba0b..4908c0d 100644 --- a/src/shared/wechat-action.ts +++ b/src/shared/wechat-action.ts @@ -2,6 +2,14 @@ export type WechatActionOrigin = | 'member_monitor' | 'scheduled_report' | 'user_tts' + /** + * 用户在界面上**手动**触发的发送(发送日报图片、给群发统计文本等)。 + * + * 立项规则:**所有发送都必须经过 `WechatActionGateway`**。手动发送在改造前是 + * 从 IPC 直连 `PersonalWechatSendService` 的,因此没有审计记录、没有幂等、 + * 也不进 Send Log —— 排查「消息到底发没发出去」时会缺一份证据。 + */ + | 'user_manual' /** Automation v1(@我生成日报)。与 `shared/automation.ts` 的 AUTOMATION_SEND_ORIGIN 一致。 */ | 'automation' | 'unknown' @@ -11,6 +19,10 @@ export type WechatActionPurpose = | 'member_left_notification' | 'scheduled_report' | 'tts_voice' + /** 用户手动发送图片(日报图片 / 群成员统计图)。 */ + | 'manual_image' + /** 用户手动发送文本。 */ + | 'manual_text' /** * Automation v1 的两个用途。 * diff --git a/tests/component/automation-rule-editor.test.tsx b/tests/component/automation-rule-editor.test.tsx new file mode 100644 index 0000000..949fdf4 --- /dev/null +++ b/tests/component/automation-rule-editor.test.tsx @@ -0,0 +1,118 @@ +import { fireEvent, render, screen } from '@testing-library/react' +import { describe, expect, it } from 'vitest' + +import { RuleEditorPanel } from '../../src/renderer/src/features/automation/RuleEditorPanel' +import { + createDefaultDailyReportRule, + type AutomationRule, + type AutomationRuleDraft +} from '../../src/shared/automation' + +/** + * 「回复前等待」是**规则自己的一项执行参数**,配置位置就是 + * 「编辑自动化 → 3 · 触发后执行(命中条件后,TraceMemo 会按下面的顺序依次执行)」。 + * + * 这一层要锁的是**接线**:初值必须来自规则本身、保存必须带回去。 + * 因为 `draftFromRule` 是逐字段搬运的,漏一个字段就会被 `normalizeRuleDraft` + * 的默认值悄悄覆盖(把用户设成 0 的等待重置回 2 秒),这种 bug 只有真渲染才发现。 + */ +const GROUPS = [{ id: '12345678@chatroom', name: '测试群' }] + +function renderEditor( + rule: AutomationRule, + onSave: (draft: AutomationRuleDraft) => void = () => {} +): void { + render( + {}} + onSave={onSave} + /> + ) +} + +const delayField = (): HTMLInputElement => + screen.getByLabelText('回复前等待秒数') as HTMLInputElement + +describe('RuleEditorPanel · 回复前等待', () => { + it('它就在「3 · 触发后执行」里,和动作顺序在同一段', () => { + renderEditor(createDefaultDailyReportRule(0)) + + expect(screen.getByText('命中条件后,TraceMemo 会按下面的顺序依次执行。')).toBeInTheDocument() + expect(delayField()).toBeInTheDocument() + expect(screen.getByText('回复前等待(秒)')).toBeInTheDocument() + }) + + it('初值来自规则本身,不是硬编码的默认值', () => { + const rule = createDefaultDailyReportRule(0) + rule.replyDelaySeconds = 7 + renderEditor(rule) + + expect(delayField()).toHaveValue(7) + }) + + it('保存时把改动带回去', () => { + const saved: AutomationRuleDraft[] = [] + renderEditor(createDefaultDailyReportRule(0), (draft) => saved.push(draft)) + + fireEvent.change(delayField(), { target: { value: '0' } }) + fireEvent.click(screen.getByRole('button', { name: '保存' })) + + expect(saved).toHaveLength(1) + expect(saved[0].replyDelaySeconds).toBe(0) + }) + + it('用户设成 0(立刻回复)时,保存不会把它顶回默认 2 秒', () => { + const saved: AutomationRuleDraft[] = [] + const rule = createDefaultDailyReportRule(0) + rule.replyDelaySeconds = 0 + renderEditor(rule, (draft) => saved.push(draft)) + + // 什么都不改,直接保存:0 必须原样保留(而不是被默认的 2 顶掉)。 + fireEvent.click(screen.getByRole('button', { name: '保存' })) + + expect(saved[0].replyDelaySeconds).toBe(0) + }) +}) + +/** + * 「所有群聊」的边界提示。 + * + * 真实能力是:同一时刻多个会话同时来消息时可能漏掉少量触发。 + * 这条限制**只挂在「所有群聊」这个上下文里** —— 首页不挂、指定群聊时不挂, + * 而且只讲用户能理解的结果,不解释表事件 / 会话回读这类实现细节。 + */ +const LIMITATION_TEXT = '当前版本在多个会话同时收到消息时,极少数自动化触发可能遗漏。' + +function renderScopeEditor(conversationIds: string[]): void { + const rule = createDefaultDailyReportRule(0) + rule.conditions.conversationIds = conversationIds + renderEditor(rule) +} + +describe('RuleEditorPanel · 生效范围的边界提示', () => { + it('一个群都没选(所有群聊)时显示提示', () => { + renderScopeEditor([]) + + // 编辑页与右侧效果预览都会显示范围描述,所以这里用 getAllByText。 + expect(screen.getAllByText('所有群聊').length).toBeGreaterThan(0) + expect(screen.getAllByText(LIMITATION_TEXT)).toHaveLength(1) + }) + + it('选了具体群聊时不显示提示', () => { + renderScopeEditor([GROUPS[0].id]) + + expect(screen.getAllByText('已选 1 个群').length).toBeGreaterThan(0) + expect(screen.queryAllByText(LIMITATION_TEXT)).toHaveLength(0) + }) + + it('提示里不出现底层实现词,也不做绝对承诺', () => { + renderScopeEditor([]) + + expect(LIMITATION_TEXT).not.toMatch(/WCDB|SessionTable|表变化|native|回读/) + expect(document.body.textContent).not.toMatch(/覆盖所有群聊|100%|不会遗漏/) + }) +}) diff --git a/tests/component/automation-workspace-status.test.tsx b/tests/component/automation-workspace-status.test.tsx new file mode 100644 index 0000000..21ae7f9 --- /dev/null +++ b/tests/component/automation-workspace-status.test.tsx @@ -0,0 +1,44 @@ +import { render, screen } from '@testing-library/react' +import { describe, expect, it } from 'vitest' + +import { ToastProvider } from '../../src/renderer/src/components/ui' +import { AutomationWorkspace } from '../../src/renderer/src/features/automation/AutomationWorkspace' + +/** + * 首页状态卡只讲用户能理解的状态。 + * + * 以前这里挂着一句「当前仅能捕获最近活跃的会话」,还用 Tooltip 解释了 + * 「微信底层只会上报有表发生了变化」——那是开发报告的内容,不该出现在普通用户界面。 + * 真实边界改挂在「规则编辑页 → 所有群聊」那个具体上下文里。 + */ +describe('Automation 首页状态卡', () => { + function renderWorkspace(): void { + render( + + + + ) + } + + it('仍然如实展示「消息监听」与发送能力两项', () => { + renderWorkspace() + + expect(screen.getByText('消息监听')).toBeInTheDocument() + expect(screen.getByText('发送能力')).toBeInTheDocument() + }) + + it('不再出现「最近活跃的会话」这类实现说明', () => { + renderWorkspace() + + expect(screen.queryByText(/最近活跃的会话/)).toBeNull() + expect(screen.queryByText(/有表发生了变化/)).toBeNull() + }) + + it('首页不出现底层实现词', () => { + renderWorkspace() + + expect(document.body.textContent).not.toMatch( + /WCDB|SessionTable|表变化|native event|回读机制|coalesce/ + ) + }) +}) diff --git a/tests/integration/automation-chain.test.ts b/tests/integration/automation-chain.test.ts index 4af976c..14b2528 100644 --- a/tests/integration/automation-chain.test.ts +++ b/tests/integration/automation-chain.test.ts @@ -131,6 +131,10 @@ function buildChain(): Chain { const store = new AutomationRuleStore({ userDataPath: () => root }) const log = new AutomationExecutionLogService({ userDataPath: () => root }) + // 回复等待归零:端到端链路要跑得快且可复现。 + // 真实默认值是 2 秒(编辑自动化 → 3 · 触发后执行),等待语义本身由 runner 单测覆盖。 + const seeded = store.listRules()[0] + if (seeded) store.updateRule(seeded.id, { ...seeded, replyDelaySeconds: 0 }) const runner = new AutomationActionRunner({ executeAction: (request) => gateway.execute(request), generateReport: async (request) => { diff --git a/tests/unit/automation-action-runner.test.ts b/tests/unit/automation-action-runner.test.ts index cbfb4c7..d334f8e 100644 --- a/tests/unit/automation-action-runner.test.ts +++ b/tests/unit/automation-action-runner.test.ts @@ -53,6 +53,14 @@ function fakeClock(step = 5): () => number { } } +/** + * 默认不真的等。 + * + * 内置规则的 `replyDelaySeconds` 默认是 2 秒,如果每条例都真 sleep, + * 这个文件会平白多出几十秒。**等待语义**由「回复前等待」那组用例显式覆盖。 + */ +const noDelay = async (): Promise => {} + function sentResult(): WechatActionResult { return { actionId: 'action-1', @@ -78,6 +86,7 @@ function failedResult(code: string, reason: string): WechatActionResult { function buildRunner(options: { executeAction?: (request: WechatActionRequest) => Promise generateReport?: () => Promise + delay?: (ms: number) => Promise }): { runner: AutomationActionRunner calls: WechatActionRequest[] @@ -87,6 +96,7 @@ function buildRunner(options: { const reports: Array<{ group: string; range?: string }> = [] const runner = new AutomationActionRunner({ now: fakeClock(), + delay: options.delay ?? noDelay, executeAction: async (request) => { calls.push(request) return options.executeAction ? options.executeAction(request) : sentResult() @@ -184,6 +194,8 @@ describe('AutomationActionRunner', () => { expect(statusOf(result, 'report')).toBe('failed') expect(statusOf(result, 'send')).toBe('skipped') expect(result.errorSummary).toBe('所选时间范围没有可总结的消息') + // 跳过必须写清「是因为前置那步挂了」,否则用户看到「已跳过」会以为规则没配好。 + expect(result.steps.find((step) => step.key === 'send')?.skipReason).toContain('前置步骤失败') // 只发过那条文字回复,图片一次都没发。 expect(calls).toHaveLength(1) expect(calls[0].content.type).toBe('text') @@ -245,16 +257,38 @@ describe('AutomationActionRunner', () => { expect(calls[0].content.type).toBe('image') }) - it('未启用生成日报时,发送图片也必须 skipped(不能凭空发图)', async () => { + it('规则没启用「发送日报图片」时,该步 skipped 且整体仍算成功(情况 A)', async () => { const { runner, calls, reports } = buildRunner({}) + const rule = makeRule((item) => { + item.actions = item.actions.filter((action) => action.type !== 'sendReportImage') + }) + const result = await run(runner, rule) + + expect(result.status).toBe('success') + expect(statusOf(result, 'reply')).toBe('success') + expect(statusOf(result, 'report')).toBe('success') + expect(statusOf(result, 'send')).toBe('skipped') + // 没配发送就只发了一条回复,不该有多余的图片请求。 + expect(calls).toHaveLength(1) + expect(calls[0].content.type).toBe('text') + expect(reports).toHaveLength(1) + }) + + it('规则要求发送图片、但手上没有图片时判失败,不允许静默跳过(情况 B)', async () => { + const { runner, calls } = buildRunner({}) + // 只留 reply + send:没有 generateReport ⇒ 永远拿不到 pngPath。 const rule = makeRule((item) => { item.actions = item.actions.filter((action) => action.type !== 'generateReport') }) const result = await run(runner, rule) expect(statusOf(result, 'report')).toBe('skipped') - expect(statusOf(result, 'send')).toBe('skipped') - expect(reports).toHaveLength(0) + expect(statusOf(result, 'send')).toBe('failed') + // 关键:整次执行必须是 failed(合并分支会把它错记成 success)。 + expect(result.status).toBe('failed') + const sendStep = result.steps.find((step) => step.key === 'send') + expect(sendStep?.error).toContain('日报图片未生成') + // 而且**绝不**退而求其次去发别的东西:只有那条文字回复出去过。 expect(calls).toHaveLength(1) expect(calls[0].content.type).toBe('text') }) @@ -287,6 +321,7 @@ describe('AutomationActionRunner', () => { const executeAction = vi.fn(async (_request: WechatActionRequest) => sentResult()) const runner = new AutomationActionRunner({ now: fakeClock(), + delay: noDelay, executeAction, generateReport: async () => ({ success: true, pngPath: '/tmp/report.png' }) }) @@ -304,4 +339,95 @@ describe('AutomationActionRunner', () => { id: 'wxid_friend' }) }) + + /** + * 回复前等待:规则一命中就秒回看起来像机器人,所以它是**规则自己的一项执行参数** + * (`rule.replyDelaySeconds`,在「编辑自动化 → 3 · 触发后执行」里配), + * 而不是全局设置 —— 不同规则可以不一样。 + * + * 这一组锁三件事:等待真的发生在回复之前、等待不计入回复步骤耗时、 + * 以及「不该等的时候一步都不等」。 + */ + it('回复等待发生在发送回复之前,顺序为 等待 → 回复 → 日报 → 图片', async () => { + const order: string[] = [] + const delay = vi.fn(async (ms: number) => { + order.push(`delay:${ms}`) + }) + const runner = new AutomationActionRunner({ + now: fakeClock(), + delay, + executeAction: async (request) => { + order.push(`send:${request.content.type}`) + return sentResult() + }, + generateReport: async () => { + order.push('report') + return { success: true, pngPath: '/tmp/report.png' } + } + }) + + await run( + runner, + makeRule((rule) => { + rule.replyDelaySeconds = 2 + }) + ) + + expect(delay).toHaveBeenCalledTimes(1) + expect(delay).toHaveBeenCalledWith(2_000) + expect(order).toEqual(['delay:2000', 'send:text', 'report', 'send:image']) + }) + + it('等待不计入「回复确认」这一步的耗时', async () => { + const { runner } = buildRunner({ delay: noDelay }) + const result = await run( + runner, + makeRule((rule) => { + rule.replyDelaySeconds = 2 + }) + ) + + // fakeClock 每调用一次进 5ms:reply 从 startedAt 到 finishedAt 只跨一次 now()。 + const replyStep = result.steps.find((step) => step.key === 'reply') + expect(replyStep?.durationMs).toBe(5) + }) + + it('等待为 0 时完全不等待;字段缺失(旧版 rules.json)按默认 2 秒', async () => { + const delay = vi.fn(noDelay) + const { runner } = buildRunner({ delay }) + + await run( + runner, + makeRule((rule) => { + rule.replyDelaySeconds = 0 + }) + ) + expect(delay).not.toHaveBeenCalled() + + // 历史规则里没有这个字段:按默认 2 秒处理,而不是凭空变成「不等待」。 + await run( + runner, + makeRule((rule) => { + delete (rule as { replyDelaySeconds?: number }).replyDelaySeconds + }) + ) + expect(delay).toHaveBeenCalledWith(2_000) + }) + + it('回复动作被关掉时不再空等', async () => { + const delay = vi.fn(noDelay) + const { runner } = buildRunner({ delay }) + + await run( + runner, + makeRule((rule) => { + rule.replyDelaySeconds = 2 + for (const action of rule.actions) { + if (action.type === 'replyText') action.enabled = false + } + }) + ) + + expect(delay).not.toHaveBeenCalled() + }) }) diff --git a/tests/unit/automation-reply-delay.test.ts b/tests/unit/automation-reply-delay.test.ts new file mode 100644 index 0000000..a796696 --- /dev/null +++ b/tests/unit/automation-reply-delay.test.ts @@ -0,0 +1,84 @@ +import { describe, expect, it } from 'vitest' + +import { + AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS, + AUTOMATION_REPLY_DELAY_MAX_SECONDS, + createDefaultDailyReportRule, + normalizeReplyDelaySeconds, + normalizeRuleDraft +} from '../../src/shared/automation' + +/** + * 「回复前等待」的口径:它是**规则自己的一项执行参数**(`replyDelaySeconds`,秒), + * 与 `cooldownSeconds` 同口径 —— 界面上是秒、落盘也是秒,只有 runner 换算成毫秒。 + * + * 核心契约:**填错了不该静默变成「不等待」**。默认 2 秒、0 是合法值、 + * 垃圾输入回落默认、超限截断。 + */ +describe('normalizeReplyDelaySeconds', () => { + it('缺省 / 空串 / null 都回落到默认 2 秒', () => { + for (const value of [undefined, null, '']) { + expect(normalizeReplyDelaySeconds(value)).toBe(AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS) + } + expect(AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS).toBe(2) + }) + + it('非法输入回落到默认值,而不是 0', () => { + for (const value of ['abc', NaN, Infinity, -Infinity, '两秒', {}]) { + expect(normalizeReplyDelaySeconds(value)).toBe(AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS) + } + }) + + it('0 与负数都收敛成「不等待」,但两者语义不同 —— 0 是显式合法值', () => { + expect(normalizeReplyDelaySeconds(0)).toBe(0) + expect(normalizeReplyDelaySeconds('0')).toBe(0) + expect(normalizeReplyDelaySeconds(-1)).toBe(0) + }) + + it('正常值取整保留,超过上限按上限截断', () => { + expect(normalizeReplyDelaySeconds(5)).toBe(5) + expect(normalizeReplyDelaySeconds('30')).toBe(30) + expect(normalizeReplyDelaySeconds(4.6)).toBe(5) + expect(normalizeReplyDelaySeconds(AUTOMATION_REPLY_DELAY_MAX_SECONDS + 1)).toBe( + AUTOMATION_REPLY_DELAY_MAX_SECONDS + ) + expect(normalizeReplyDelaySeconds(99_999)).toBe(AUTOMATION_REPLY_DELAY_MAX_SECONDS) + }) +}) + +describe('规则里的回复前等待', () => { + it('内置「@我生成日报」默认 2 秒', () => { + expect(createDefaultDailyReportRule(0).replyDelaySeconds).toBe( + AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS + ) + }) + + it('草稿归一化会带上这个字段(否则编辑页拿不到值)', () => { + const draft = normalizeRuleDraft({ + name: '测试', + replyDelaySeconds: 7, + actions: [{ type: 'replyText', enabled: true, text: '收到' }] + }) + expect(draft.replyDelaySeconds).toBe(7) + }) + + it('旧版 rules.json 里的规则没有这个字段 → 读出来是默认 2 秒,不是 undefined', () => { + const draft = normalizeRuleDraft({ + name: '旧规则', + actions: [{ type: 'replyText', enabled: true, text: '收到' }] + }) + expect(draft.replyDelaySeconds).toBe(AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS) + }) + + it('草稿里的脏值同样被收敛', () => { + expect( + normalizeRuleDraft({ name: 'x', replyDelaySeconds: -3, actions: [] }).replyDelaySeconds + ).toBe(0) + expect( + normalizeRuleDraft({ name: 'x', replyDelaySeconds: 'abc', actions: [] }).replyDelaySeconds + ).toBe(AUTOMATION_REPLY_DELAY_DEFAULT_SECONDS) + expect( + normalizeRuleDraft({ name: 'x', replyDelaySeconds: 1e9, actions: [] }).replyDelaySeconds + ).toBe(AUTOMATION_REPLY_DELAY_MAX_SECONDS) + }) +}) diff --git a/tests/unit/automation-service.test.ts b/tests/unit/automation-service.test.ts index 1dbb27c..8c2be30 100644 --- a/tests/unit/automation-service.test.ts +++ b/tests/unit/automation-service.test.ts @@ -100,7 +100,18 @@ interface Harness { setRunResult: (value: AutomationRunResult) => void } -function buildHarness(options: { now?: () => number; capability?: PersonalWechatSendCapability } = {}): Harness { +function buildHarness( + options: { + now?: () => number + capability?: PersonalWechatSendCapability + /** 执行期间的挂起点:resolve 之前 `runner.run` 不会返回(用来模拟长任务)。 */ + holdExecution?: () => Promise + /** 每次真正进入 runner 时通知测试,便于确定时序。 */ + onRunStart?: () => void + /** gate 表容量上限,用来在单测里触发清理路径。 */ + gateMaxEntries?: number + } = {} +): Harness { const store = new AutomationRuleStore() const log = new AutomationExecutionLogService() const runs: AutomationRunInput[] = [] @@ -120,6 +131,8 @@ function buildHarness(options: { now?: () => number; capability?: PersonalWechat const spyRunner = { run: async (input: AutomationRunInput): Promise => { runs.push(input) + options.onRunStart?.() + await options.holdExecution?.() return runResult } } as unknown as AutomationActionRunner @@ -135,7 +148,8 @@ function buildHarness(options: { now?: () => number; capability?: PersonalWechat runner: spyRunner, getCapability: async () => currentCapability, isListening: () => true, - ...(options.now ? { now: options.now } : {}) + ...(options.now ? { now: options.now } : {}), + ...(options.gateMaxEntries !== undefined ? { gateMaxEntries: options.gateMaxEntries } : {}) }) return { @@ -187,13 +201,17 @@ describe('AutomationService', () => { expect(records[0].executionId).toBe(harness.runs[0].executionId) }) - it('同一条消息被重复投递时只执行一次', async () => { + it('同一条消息被重复投递时只执行一次,且**不产生任何额外记录**', async () => { const harness = buildHarness() + // 关掉 cooldown,单独验证「同消息幂等」这条闸 —— 否则第二次会先被 gate 拦下。 + const seeded = harness.store.listRules()[0] + harness.store.updateRule(seeded.id, { ...seeded, cooldownSeconds: 0 }) + await harness.service.handleMessage(message()) await harness.service.handleMessage(message()) - await harness.service.handleMessage(message({ localId: '100' })) expect(harness.runs).toHaveLength(1) + // 重复投递是「这条消息已经处理过」,不是「跳过了一次执行」——用户日志里不该出现。 expect(harness.log.list()).toHaveLength(1) }) @@ -205,7 +223,7 @@ describe('AutomationService', () => { expect(harness.runs).toHaveLength(2) }) - it('cooldown 内不重复触发', async () => { + it('cooldown 内不重复触发,且对用户完全静默', async () => { let now = 1_700_000_000_000 const harness = buildHarness({ now: () => now }) @@ -213,7 +231,12 @@ describe('AutomationService', () => { now += 30_000 await harness.service.handleMessage(message({ localId: '2' })) + // 只执行一次 —— 这条是硬要求。 expect(harness.runs).toHaveLength(1) + // 但第二次**不留任何痕迹**:没有 execution、没有「已跳过」、没有「冷却中」。 + expect(harness.log.list()).toHaveLength(1) + // 只进诊断计数(不上 UI)。 + expect(harness.service.getBlockedMessageCount()).toBe(1) }) it('cooldown 过后可以再次触发', async () => { @@ -370,6 +393,16 @@ describe('AutomationService', () => { }) }) + it('把命中时的规则整条交给 Runner(含回复前等待)', async () => { + const harness = buildHarness() + await harness.service.handleMessage(message()) + + expect(harness.runs).toHaveLength(1) + // 「回复前等待」是规则自己的字段,编排层只负责把规则原样传下去 —— + // 不该再有一份全局设置参与,否则同一字段有两个来源就必然对不上。 + expect(harness.runs[0].rule.replyDelaySeconds).toBe(2) + }) + it('handleMessage 内部异常不会向外抛', async () => { const harness = buildHarness() const brokenStore = { @@ -385,3 +418,249 @@ describe('AutomationService', () => { void harness }) }) + +/** + * Automation rule conversation gate —— 规则 × 会话维度的阻塞窗口。 + * + * ``` + * blocked = inFlight || now < triggeredAt + cooldown // 粒度:ruleId + conversationId + * ``` + * + * 窗口内的消息:不匹配、不执行、不回复、**不写用户执行日志**,只进诊断计数。 + * `triggeredAt` 是**第一次触发**的时间;执行完成不重置它。 + */ +describe('AutomationService · rule conversation gate', () => { + const T0 = 1_700_000_000_000 + + // 这个 describe 在文件后段,**不会**继承上一个 describe 的 beforeEach; + // 而规则是写盘的(同一个 userData root),不清理就会互相污染 + // (Test 6 建的第二条规则会漏给 Test 7/8)。 + beforeEach(() => { + fs.removeSync(path.join(root, 'automation')) + }) + + it('Test 1 · 0s 触发,2/20/59s 全忽略,61s 重新触发(executions = 2 而不是 5)', async () => { + let now = T0 + const harness = buildHarness({ now: () => now }) + + await harness.service.handleMessage(message({ localId: '1' })) + for (const offset of [2_000, 20_000, 59_000]) { + now = T0 + offset + await harness.service.handleMessage(message({ localId: `t${offset}` })) + } + now = T0 + 61_000 + await harness.service.handleMessage(message({ localId: 'e' })) + + expect(harness.runs).toHaveLength(2) + expect(harness.log.list()).toHaveLength(2) + expect(harness.service.getBlockedMessageCount()).toBe(3) + }) + + it('Test 2 · 短任务(8s 就完成)仍然阻塞到 cooldown 到期', async () => { + let now = T0 + const harness = buildHarness({ now: () => now }) + + await harness.service.handleMessage(message({ localId: '1' })) // 0s 触发,立即完成 + now = T0 + 8_000 + now = T0 + 30_000 // 任务早完成了,但 cooldown 还没到 + await harness.service.handleMessage(message({ localId: '2' })) + expect(harness.runs).toHaveLength(1) + + now = T0 + 61_000 + await harness.service.handleMessage(message({ localId: '3' })) + expect(harness.runs).toHaveLength(2) + }) + + it('Test 3 · 长任务(超过 cooldown)一直阻塞到执行结束:解锁 = max(完成, 触发+cooldown)', async () => { + let now = T0 + let release!: () => void + const held = new Promise((resolve) => { + release = resolve + }) + let markStarted!: () => void + const started = new Promise((resolve) => { + markStarted = resolve + }) + const harness = buildHarness({ + now: () => now, + holdExecution: () => held, + onRunStart: markStarted + }) + + const firstRun = harness.service.handleMessage(message({ localId: '1' })) + await started + + // 61s:cooldown 早就过期了,但第一条执行还在跑 ⇒ 仍然忽略 + now = T0 + 61_000 + await harness.service.handleMessage(message({ localId: '2' })) + expect(harness.runs).toHaveLength(1) + + // 75s:第一条终于结束 + now = T0 + 75_000 + release() + await firstRun + + // 76s:这才重新允许 —— **不是**「完成后再等 60 秒」 + now = T0 + 76_000 + await harness.service.handleMessage(message({ localId: '3' })) + expect(harness.runs).toHaveLength(2) + }) + + it('Test 4 · 同群不同发送者也忽略(gate 不是 sender 级)', async () => { + let now = T0 + const harness = buildHarness({ now: () => now }) + + await harness.service.handleMessage(message({ localId: '1' })) + now = T0 + 20_000 + await harness.service.handleMessage( + message({ localId: '2', senderId: 'wxid_other', senderNickname: '李四' }) + ) + + expect(harness.runs).toHaveLength(1) + }) + + it('Test 5 · 同一条规则在另一个群正常执行(gate 是 ruleId + conversationId)', async () => { + let now = T0 + const harness = buildHarness({ now: () => now }) + + await harness.service.handleMessage(message({ localId: '1' })) + now = T0 + 10_000 + await harness.service.handleMessage(message({ localId: '2', sessionId: OTHER_GROUP })) + + expect(harness.runs).toHaveLength(2) + }) + + it('Test 6 · 同一群里的另一条规则不受影响(不是会话全局锁)', async () => { + let now = T0 + const harness = buildHarness({ now: () => now }) + const seeded = harness.store.listRules()[0] + // 第二条规则同样命中这条消息,但**没有冷却** —— 用它证明 gate 按 ruleId 隔离。 + harness.store.createRule({ + name: '帮助', + enabled: true, + scope: 'group', + conditions: { + requireMentionMe: true, + keyword: '日报', + keywordMatchMode: 'contains', + conversationIds: [], + ignoreSelf: true + }, + actions: [{ type: 'replyText', enabled: true, text: '好的' }], + cooldownSeconds: 0, + replyDelaySeconds: 0 + }) + + await harness.service.handleMessage(message({ localId: '1' })) + expect(harness.runs).toHaveLength(2) + + now = T0 + 10_000 + await harness.service.handleMessage(message({ localId: '2' })) + // 规则 A 被 gate 挡下,规则 B 照常执行。 + expect(harness.runs).toHaveLength(3) + expect(harness.runs[2].rule.id).not.toBe(seeded.id) + }) + + it('Test 7 · 第一条失败后仍然走完整 cooldown(避免故障期疯狂重试)', async () => { + let now = T0 + const harness = buildHarness({ now: () => now }) + harness.setRunResult({ steps: [], status: 'failed', errorSummary: 'AI 服务不可达' }) + + await harness.service.handleMessage(message({ localId: '1' })) + expect(harness.runs).toHaveLength(1) + + now = T0 + 5_000 + await harness.service.handleMessage(message({ localId: '2' })) + expect(harness.runs).toHaveLength(1) + + now = T0 + 61_000 + await harness.service.handleMessage(message({ localId: '3' })) + expect(harness.runs).toHaveLength(2) + }) + + it('Test 8 · 60 秒内 10 条符合条件的消息,用户执行日志只留 1 条', async () => { + let now = T0 + const harness = buildHarness({ now: () => now }) + + await harness.service.handleMessage(message({ localId: '1' })) + for (let index = 2; index <= 11; index += 1) { + now = T0 + index * 5_000 // 5s..55s,全部落在 60s 窗口内 + await harness.service.handleMessage(message({ localId: String(index) })) + } + + expect(harness.runs).toHaveLength(1) + const records = harness.log.list() + expect(records).toHaveLength(1) + expect(records[0].status).toBe('success') + // 不许出现「冷却中 / 剩余 N 秒」这类记录。 + expect(JSON.stringify(records)).not.toContain('冷却') + expect(harness.service.getBlockedMessageCount()).toBe(10) + }) +}) + +/** + * gate 表清理的硬不变量:**正在执行(`inFlight === true`)的 gate 永不被淘汰**。 + * + * 淘汰它 = 把一条正在跑的规则提前放开,立刻会产生第二次并发执行。 + * 这里用极小的上限(1)逼出清理路径。 + */ +describe('AutomationService · gate 缓存清理', () => { + const T0 = 1_700_000_000_000 + const GROUP_A = 'aaa@chatroom' + const GROUP_B = 'bbb@chatroom' + const GROUP_C = 'ccc@chatroom' + + beforeEach(() => { + fs.removeSync(path.join(root, 'automation')) + }) + + it('清理只删「inFlight=false 且 cooldown 已过期」的条目,正在执行的一条必须留着', async () => { + let now = T0 + let release!: () => void + const held = new Promise((resolve) => { + release = resolve + }) + let markStarted!: () => void + const started = new Promise((resolve) => { + markStarted = resolve + }) + // 只把**第一次**执行挂住(也就是 A 群那条);B/C 的执行立即完成。 + let runCount = 0 + const harness = buildHarness({ + now: () => now, + gateMaxEntries: 1, + holdExecution: () => { + runCount += 1 + return runCount === 1 ? held : Promise.resolve() + }, + onRunStart: markStarted + }) + // cooldown = 0:这样「执行完成」的 gate 立刻就算「已解锁」,可以被清理。 + const seeded = harness.store.listRules()[0] + harness.store.updateRule(seeded.id, { ...seeded, cooldownSeconds: 0 }) + + // A 群:开始执行并**挂住**(gate A 处于 inFlight) + const runningA = harness.service.handleMessage( + message({ localId: 'a1', sessionId: GROUP_A }) + ) + await started + expect(harness.runs).toHaveLength(1) + + // B 群:执行一次并正常结束(gate B 变成「已解锁」,是本次清理的目标) + now = T0 + 1_000 + await harness.service.handleMessage(message({ localId: 'b1', sessionId: GROUP_B })) + expect(harness.runs).toHaveLength(2) + + // C 群:claim 一条新 gate ⇒ 表 size 超过上限 ⇒ 触发 evictGatesIfNeeded + now = T0 + 2_000 + await harness.service.handleMessage(message({ localId: 'c1', sessionId: GROUP_C })) + + // 关键不变量:A 的执行还在跑,往 A 发消息必须**仍然被挡住**。 + now = T0 + 3_000 + await harness.service.handleMessage(message({ localId: 'a2', sessionId: GROUP_A })) + expect(harness.runs.filter((run) => run.conversationId === GROUP_A)).toHaveLength(1) + + release() + await runningA + }) +}) diff --git a/tests/unit/message-listener-service.test.ts b/tests/unit/message-listener-service.test.ts index bfd71ed..dcb7044 100644 --- a/tests/unit/message-listener-service.test.ts +++ b/tests/unit/message-listener-service.test.ts @@ -6,6 +6,7 @@ import { extractMentionTargets, type NormalizedIncomingMessage } from '../../src/main/services/message-listener-service' +import { parseMonitorEvent } from '../../src/main/wcdb4-client' import type { Wcdb4Client, Wcdb4Message } from '../../src/main/wcdb4-client' /** @@ -344,3 +345,166 @@ describe('Observation Mode 约束', () => { expect(received).toHaveLength(0) }) }) + +/** + * Native Monitor Event v2:pipe payload 里带上**发生变化的会话**。 + * + * 这一组锁的是「不再依赖 `getSessions()[0]`」这件事,以及两个必须保留的行为: + * legacy(v1 runtime)fallback 和同消息 dedup。 + */ +describe('parseMonitorEvent(协议解析)', () => { + it('v2 + message_change + session_id ⇒ precise', () => { + const event = parseMonitorEvent( + JSON.stringify({ + version: 2, + kind: 'message_change', + session_id: 'group_A@chatroom', + db: 'message_0.db', + table: 'message', + action: 'update', + observed_at_ms: 1_700_000_000_000 + }) + ) + expect(event.protocol).toBe(2) + expect(event.sessionId).toBe('group_A@chatroom') + expect(event.table).toBe('message') + expect(event.observedAtMs).toBe(1_700_000_000_000) + }) + + it('v1 payload 仍然解析成 legacy(不带会话)', () => { + const event = parseMonitorEvent( + JSON.stringify({ db: 'session.db', table: 'Session', action: 'update' }) + ) + expect(event.protocol).toBe(1) + expect(event.sessionId).toBeUndefined() + expect(event.action).toBe('update') + expect(event.raw).toContain('session.db') + }) + + it('v2 但缺 session_id / kind 不对 ⇒ 一律降级成 legacy,绝不猜会话', () => { + for (const payload of [ + { version: 2, kind: 'message_change' }, + { version: 2, kind: 'message_change', session_id: ' ' }, + { version: 2, session_id: 'group_A@chatroom' }, + { version: 1, kind: 'message_change', session_id: 'group_A@chatroom' } + ]) { + expect(parseMonitorEvent(JSON.stringify(payload)).protocol).toBe(1) + expect(parseMonitorEvent(JSON.stringify(payload)).sessionId).toBeUndefined() + } + }) + + it('不是 JSON 也不崩', () => { + const event = parseMonitorEvent('not-json') + expect(event.protocol).toBe(1) + expect(event.action).toBe('update') + }) +}) + +describe('MessageListener · precise session path', () => { + const v2Event = (sessionId: string) => + parseMonitorEvent( + JSON.stringify({ version: 2, kind: 'message_change', session_id: sessionId, table: 'message' }) + ) + const v1Event = () => + parseMonitorEvent(JSON.stringify({ db: 'session.db', table: 'Session', action: 'update' })) + + /** 按会话返回消息,并记录每个会话被回读了几次。 */ + function multiSessionClient(bySession: Record): { + client: Wcdb4Client + calls: string[] + } { + const calls: string[] = [] + const client = { + // legacy fallback 用的「最近活跃会话」,故意与 precise 的会话不同。 + getSessions: vi.fn(() => [{ username: 'recent@chatroom' }]), + getMessagesAsync: vi.fn(async (sessionId: string) => { + calls.push(sessionId) + return bySession[sessionId] ?? [] + }) + } as unknown as Wcdb4Client + return { client, calls } + } + + it('v2 事件直接按 sessionId 回读,不再碰 getSessions()[0]', async () => { + const { client, calls } = multiSessionClient({ + 'A@chatroom': [makeMessage({ mesLocalID: '1' })] + }) + const listener = new MessageListenerService(client) + collect(listener) + + listener.handleNativeChange(v2Event('A@chatroom')) + await flush() + + expect(calls).toEqual(['A@chatroom']) + expect(received.map((message) => message.sessionId)).toEqual(['A@chatroom']) + }) + + it('120ms 内 A/B/A/C 会回读 A、B、C 各一次(不是「最后一个 wins」)', async () => { + const { client, calls } = multiSessionClient({ + 'A@chatroom': [makeMessage({ mesLocalID: '1' })], + 'B@chatroom': [makeMessage({ mesLocalID: '2' })], + 'C@chatroom': [makeMessage({ mesLocalID: '3' })] + }) + const listener = new MessageListenerService(client) + collect(listener) + + listener.handleNativeChange(v2Event('A@chatroom')) + listener.handleNativeChange(v2Event('B@chatroom')) + listener.handleNativeChange(v2Event('A@chatroom')) + listener.handleNativeChange(v2Event('C@chatroom')) + await flush() + + expect(calls.sort()).toEqual(['A@chatroom', 'B@chatroom', 'C@chatroom']) + expect(received.map((message) => message.sessionId).sort()).toEqual([ + 'A@chatroom', + 'B@chatroom', + 'C@chatroom' + ]) + }) + + it('一批里既有 v2 又有 v1 时,仍然按 v2 的精确会话回读', async () => { + const { client, calls } = multiSessionClient({ + 'A@chatroom': [makeMessage({ mesLocalID: '1' })], + 'recent@chatroom': [makeMessage({ mesLocalID: '9' })] + }) + const listener = new MessageListenerService(client) + collect(listener) + + listener.handleNativeChange(v2Event('A@chatroom')) + listener.handleNativeChange(v1Event()) + await flush() + + expect(calls).toEqual(['A@chatroom']) + }) + + it('legacy(v1 runtime)仍然退回最近活跃会话,行为不变', async () => { + const { client, calls } = multiSessionClient({ + 'recent@chatroom': [makeMessage({ mesLocalID: '1' })] + }) + const listener = new MessageListenerService(client) + collect(listener) + + listener.handleNativeChange(v1Event()) + await flush() + + expect(calls).toEqual(['recent@chatroom']) + expect(received).toHaveLength(1) + }) + + it('native 重复发同一条消息的 v2 事件,只会投递一次', async () => { + const { client, calls } = multiSessionClient({ + 'A@chatroom': [makeMessage({ mesLocalID: '42' })] + }) + const listener = new MessageListenerService(client) + collect(listener) + + listener.handleNativeChange(v2Event('A@chatroom')) + await flush() + listener.handleNativeChange(v2Event('A@chatroom')) + await flush() + + // 回读了两次(两次事件各一个窗口),但同一条消息只投递一次。 + expect(calls).toEqual(['A@chatroom', 'A@chatroom']) + expect(received).toHaveLength(1) + }) +}) diff --git a/tests/unit/wechat-action-gateway.test.ts b/tests/unit/wechat-action-gateway.test.ts index 2089eb4..7a56f2c 100644 --- a/tests/unit/wechat-action-gateway.test.ts +++ b/tests/unit/wechat-action-gateway.test.ts @@ -18,10 +18,11 @@ vi.mock('../../src/main/services/personal-wechat-send-service', () => ({ import { AUTOMATION_SEND_INTERVAL_MS, - WechatActionGateway + WechatActionGateway, + toPersonalWechatSendResult } from '../../src/main/services/wechat-action-gateway' import type { PersonalWechatSendCapability } from '../../src/shared/personal-wechat' -import type { WechatActionResult } from '../../src/shared/wechat-action' +import type { WechatActionRequest, WechatActionResult } from '../../src/shared/wechat-action' const readyCapability: PersonalWechatSendCapability = { supported: true, @@ -375,4 +376,132 @@ describe('WechatActionGateway', () => { expect(Date.now()).toBe(startedAt) vi.useRealTimers() }) + + /** + * 项目规则:**所有发送都走 WechatActionGateway**。 + * + * 手动发送(`wechat-personal:send`)改造前是直连 `PersonalWechatSendService` 的, + * 于是没有审计、没有幂等、也不进 Send Log。这一组锁住改造后的语义。 + */ + describe('用户手动发送(triggerType: user)', () => { + function manualAction( + gateway: WechatActionGateway, + overrides: Partial = {} + ): Promise { + return gateway.execute({ + origin: 'user_manual', + purpose: 'manual_image', + triggerType: 'user', + recipient: { type: 'group', id: 'room@chatroom' }, + content: { type: 'image', path: '/tmp/report.png' }, + ...overrides + } as WechatActionRequest) + } + + it('放行、发送一次,并写入审计记录', async () => { + const userData = mkdtempSync(join(tmpdir(), 'tracememo-wechat-action-')) + directories.push(userData) + const gateway = new WechatActionGateway({ getUserDataPath: () => userData }) + + const result = await manualAction(gateway) + + expect(result).toMatchObject({ status: 'sent', decision: 'allow' }) + expect(mocks.sender.send).toHaveBeenCalledWith({ + type: 'image', + to: 'room@chatroom', + isGroup: true, + filePath: '/tmp/report.png' + }) + const audit = readJsonSync(join(userData, 'actions', 'wechat-actions.json')) + expect(audit).toEqual([ + expect.objectContaining({ + origin: 'user_manual', + purpose: 'manual_image', + triggerType: 'user', + recipientId: 'room@chatroom', + decision: 'allow', + sendStatus: 'sent' + }) + ]) + }) + + it('不受 automation 的 purpose allowlist 限制,也不吃 3s 节流', async () => { + vi.useFakeTimers() + vi.setSystemTime(new Date('2026-09-21T01:00:00.000Z')) + const gateway = createGateway() + + const first = await manualAction(gateway) + const second = await manualAction(gateway) + + expect(first.status).toBe('sent') + expect(second.status).toBe('sent') + // 节流只对 triggerType === 'automation' 生效:两次点击必须立刻都发出去。 + expect(mocks.sender.send).toHaveBeenCalledTimes(2) + vi.useRealTimers() + }) + + it('空文本由规范化拦下,不落到发送层', async () => { + const gateway = createGateway() + const result = await manualAction(gateway, { + purpose: 'manual_text', + content: { type: 'text', text: ' ' } + }) + + expect(result.status).toBe('blocked') + expect(result.errorCode).toBe('INVALID_REQUEST') + expect(mocks.sender.send).not.toHaveBeenCalled() + }) + }) + + describe('toPersonalWechatSendResult', () => { + const fallback = { canSend: true } as unknown as Parameters< + typeof toPersonalWechatSendResult + >[1] + + it('优先返回底层 sendResult(它带着真实的 status)', () => { + const underlying = { success: true, status: { canSend: true, marker: 'real' } } + expect( + toPersonalWechatSendResult( + { + actionId: 'a1', + status: 'sent', + decision: 'allow', + startedAt: '2026-09-21T01:00:00.000Z', + finishedAt: '2026-09-21T01:00:01.000Z', + sendResult: underlying + }, + fallback + ) + ).toBe(underlying) + }) + + it('拿不到 sendResult 时按 action.status 合成', () => { + const sent = toPersonalWechatSendResult( + { + actionId: 'a1', + status: 'sent', + decision: 'allow', + startedAt: '2026-09-21T01:00:00.000Z', + finishedAt: '2026-09-21T01:00:01.000Z' + }, + fallback + ) + expect(sent).toEqual({ success: true, status: fallback }) + + const blocked = toPersonalWechatSendResult( + { + actionId: 'a2', + status: 'blocked', + decision: 'block', + errorCode: 'ACTION_NOT_ALLOWED', + reason: '自动化动作不允许执行', + startedAt: '2026-09-21T01:00:00.000Z', + finishedAt: '2026-09-21T01:00:01.000Z' + }, + fallback + ) + expect(blocked.success).toBe(false) + expect(blocked.error).toBe('自动化动作不允许执行') + }) + }) })