mirror of
https://wget.la/https://github.com/Wxw-Gu/WechatExplorer
synced 2026-10-03 18:33:14 +08:00
feat: 精确会话自动化规则
This commit is contained in:
Binary file not shown.
+41
-10
@@ -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)
|
||||
|
||||
@@ -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<AgentGroupReportResult>
|
||||
executeAction?: (request: WechatActionRequest) => Promise<WechatActionResult>
|
||||
now?: () => number
|
||||
/** 延迟实现。默认真 sleep;单测注入即时 resolve 的假实现,避免真的等 2 秒。 */
|
||||
delay?: (ms: number) => Promise<void>
|
||||
}
|
||||
|
||||
/** 策略层的错误码 → 用户可读短句。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<AgentGroupReportResult>
|
||||
private readonly executeAction: (request: WechatActionRequest) => Promise<WechatActionResult>
|
||||
private readonly now: () => number
|
||||
private readonly delay: (ms: number) => Promise<void>
|
||||
|
||||
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<void>((resolve) => setTimeout(resolve, ms)))
|
||||
}
|
||||
|
||||
async run(input: AutomationRunInput): Promise<AutomationRunResult> {
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<AutomationExecution>
|
||||
@@ -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
|
||||
|
||||
@@ -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<string, ClaimEntry>()
|
||||
/** `${ruleId}:${conversationId}` → 上次真正开始执行的毫秒时间戳。 */
|
||||
private readonly cooldowns = new Map<string, number>()
|
||||
/**
|
||||
* Automation rule ↔ conversation gate。
|
||||
*
|
||||
* 粒度是 **`ruleId + conversationId`**(不是 sender,也不是整个会话):
|
||||
* 群 A 的「@我生成日报」被触发后,群 B 的同一条规则、或同群里的**另一条规则**都不受影响。
|
||||
*/
|
||||
private readonly gates = new Map<string, ConversationRuleGate>()
|
||||
|
||||
/** 仅用于诊断统计(`[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<ReturnType<AutomationActionRunner['run']>>
|
||||
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()
|
||||
|
||||
@@ -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<string, number>()
|
||||
/**
|
||||
* 本批 coalesce 窗口内**出现过精确会话**的事件收集到的 session 集合。
|
||||
*
|
||||
* v2 事件(Native Monitor Event v2)会带上发生变化的 sessionId,这里用 **Set** 收集 ——
|
||||
* 120ms 内收到 `A B A C B` 必须回读 A/B/C 三个,绝不能「最后一个 wins」。
|
||||
*/
|
||||
private readonly pendingSessions = new Set<string>()
|
||||
private coalesceTimer: ReturnType<typeof setTimeout> | 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<void> {
|
||||
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<number> {
|
||||
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 {
|
||||
|
||||
@@ -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<string, string | number | boolean | null>): 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<string, unknown>,
|
||||
@@ -939,15 +977,41 @@ async function requestWindowsHook(
|
||||
host = windowsHookHost()
|
||||
): Promise<WindowsHookResponse> {
|
||||
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) {
|
||||
|
||||
@@ -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<string, unknown>)) {
|
||||
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<WechatActionAuditRecord, 'contentPreview' | 'contentHash'> {
|
||||
|
||||
@@ -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<string, unknown> = {}
|
||||
try {
|
||||
const value = JSON.parse(raw) as unknown
|
||||
if (value && typeof value === 'object' && !Array.isArray(value)) {
|
||||
parsed = value as Record<string, unknown>
|
||||
}
|
||||
} 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<typeof setTimeout> | null = null
|
||||
private monitorReconnectTimer: ReturnType<typeof setTimeout> | null = null
|
||||
private monitorPipePath = ''
|
||||
@@ -773,7 +846,7 @@ export class Wcdb4Client {
|
||||
})
|
||||
}
|
||||
|
||||
async startMonitor(callback: (type: string, json: string) => void): Promise<boolean> {
|
||||
async startMonitor(callback: (event: Wcdb4MonitorEvent) => void): Promise<boolean> {
|
||||
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 {
|
||||
|
||||
@@ -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({
|
||||
<div className="automation-status-card">
|
||||
<span className="automation-status-label">消息监听</span>
|
||||
<span className={`automation-status-value ${status.listening ? 'ok' : 'off'}`}>
|
||||
{status.listening ? '正在监听新消息' : '未在监听'}
|
||||
{status.listening ? '运行中' : '未在监听'}
|
||||
</span>
|
||||
{status.listeningDegraded ? (
|
||||
<Tooltip>
|
||||
<TooltipTrigger asChild>
|
||||
<small className="automation-status-note" tabIndex={0}>
|
||||
当前仅能捕获最近活跃的会话
|
||||
</small>
|
||||
</TooltipTrigger>
|
||||
<TooltipContent side="bottom">
|
||||
微信底层只会上报「有表发生了变化」,不带是哪条会话。TraceMemo
|
||||
目前据此回读最近活跃的会话,因此同一时刻多个群同时来消息时,只会处理其中的一个。
|
||||
</TooltipContent>
|
||||
</Tooltip>
|
||||
) : null}
|
||||
</div>
|
||||
|
||||
<div className="automation-status-card">
|
||||
|
||||
@@ -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({
|
||||
<dt>结果</dt>
|
||||
<dd>
|
||||
<span className={`automation-status-chip ${execution.status}`}>
|
||||
{execution.status === 'success' ? '成功' : '失败'}
|
||||
{AUTOMATION_EXECUTION_STATUS_LABELS[execution.status]}
|
||||
</span>
|
||||
</dd>
|
||||
</div>
|
||||
|
||||
@@ -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({
|
||||
</span>
|
||||
<span role="cell">
|
||||
<span className={`automation-status-chip ${execution.status}`}>
|
||||
{execution.status === 'success' ? '成功' : '失败'}
|
||||
{AUTOMATION_EXECUTION_STATUS_LABELS[execution.status]}
|
||||
</span>
|
||||
</span>
|
||||
<span role="cell" className="automation-log-duration">
|
||||
|
||||
@@ -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({
|
||||
<div className="automation-section-heading">
|
||||
<h3>2 · 在哪些聊天生效</h3>
|
||||
<span className="automation-section-note">
|
||||
{draft.conditions.conversationIds.length
|
||||
{hasSelectedGroups
|
||||
? `已选 ${draft.conditions.conversationIds.length} 个群`
|
||||
: '未选择时对所有群聊生效'}
|
||||
: '所有群聊'}
|
||||
</span>
|
||||
</div>
|
||||
<Input
|
||||
@@ -261,6 +275,14 @@ export function RuleEditorPanel({
|
||||
})
|
||||
)}
|
||||
</div>
|
||||
{hasSelectedGroups ? null : (
|
||||
// 只在「所有群聊」这个上下文里说明真实边界。
|
||||
// 不解释底层原因(表事件 / 会话回读),也不写「100% 覆盖」这种绝对承诺。
|
||||
<p className="automation-section-hint">
|
||||
<span aria-hidden="true">ⓘ</span>
|
||||
当前版本在多个会话同时收到消息时,极少数自动化触发可能遗漏。
|
||||
</p>
|
||||
)}
|
||||
</section>
|
||||
|
||||
<section className="automation-section">
|
||||
@@ -268,6 +290,30 @@ export function RuleEditorPanel({
|
||||
<h3>3 · 触发后执行</h3>
|
||||
<span className="automation-section-note">按下列顺序执行</span>
|
||||
</div>
|
||||
{/* 第一步不是动作,而是「等多久」—— 所以放在动作列表之前,不混进 ACTION_META。 */}
|
||||
<div className="automation-inline-row">
|
||||
<div>
|
||||
<span className="automation-field-label">回复前等待(秒)</span>
|
||||
<small>命中后先等这么久再回复,避免秒回显得像机器人;0 = 立刻回复</small>
|
||||
</div>
|
||||
<Input
|
||||
type="number"
|
||||
min={0}
|
||||
max={AUTOMATION_REPLY_DELAY_MAX_SECONDS}
|
||||
value={String(draft.replyDelaySeconds)}
|
||||
onChange={(event) => {
|
||||
const next = Number(event.target.value)
|
||||
setDraft((current) => ({
|
||||
...current,
|
||||
replyDelaySeconds: normalizeReplyDelaySeconds(
|
||||
Number.isFinite(next) ? next : 0
|
||||
)
|
||||
}))
|
||||
}}
|
||||
className="automation-number-input"
|
||||
aria-label="回复前等待秒数"
|
||||
/>
|
||||
</div>
|
||||
{ACTION_META.map((meta) => {
|
||||
const action = actionOf(meta.type)
|
||||
return (
|
||||
@@ -304,7 +350,10 @@ export function RuleEditorPanel({
|
||||
<div className="automation-inline-row">
|
||||
<div>
|
||||
<span className="automation-field-label">触发间隔(秒)</span>
|
||||
<small>同一个群里,两次触发之间至少间隔这么久,避免刷屏</small>
|
||||
{/* 这句就是本产品的门语义定义,不要写 debounce / mutex / in-flight 这类技术词。 */}
|
||||
<small>
|
||||
触发后,在当前任务执行期间及设定间隔内,不再处理本规则在该会话中的新消息。
|
||||
</small>
|
||||
</div>
|
||||
<Input
|
||||
type="number"
|
||||
|
||||
@@ -50,9 +50,13 @@ export function stepStatusTone(status: AutomationStepStatus): string {
|
||||
* `skipped` 步骤的说明。
|
||||
*
|
||||
* 用户最容易困惑的就是「为什么这一步没跑」—— 必须区分
|
||||
* 「规则本来就没配这个动作」和「上一步挂了所以跳过」。
|
||||
* 「规则本来就没配这个动作」「上一步挂了所以跳过」「被规则自身的冷却/去重拦下」。
|
||||
*
|
||||
* **优先用服务端给的 `skipReason`**:它知道确切原因(例如「前置步骤失败(日报未生成)」),
|
||||
* 下面这套推断只是兜底,不许在这里重新发明原因。
|
||||
*/
|
||||
export function describeSkippedStep(step: AutomationStep, steps: AutomationStep[]): string {
|
||||
if (step.skipReason) return step.skipReason
|
||||
const index = steps.findIndex((item) => 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 '未执行'
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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<AutomationStepKey, string> = {
|
||||
send: '发送日报图片'
|
||||
}
|
||||
|
||||
/**
|
||||
* 执行结果的展示文案。
|
||||
*
|
||||
* `skipped` 是**中性**状态(命中了但被规则自身的冷却 / 去重拦下),
|
||||
* 不是失败 —— 界面不能用红色渲染它。
|
||||
*/
|
||||
export const AUTOMATION_EXECUTION_STATUS_LABELS: Record<AutomationExecutionStatus, string> = {
|
||||
running: '执行中',
|
||||
success: '成功',
|
||||
failed: '失败'
|
||||
}
|
||||
|
||||
export const KEYWORD_MATCH_MODE_LABELS: Record<KeywordMatchMode, string> = {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 的两个用途。
|
||||
*
|
||||
|
||||
@@ -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(
|
||||
<RuleEditorPanel
|
||||
mode="edit"
|
||||
rule={rule}
|
||||
groups={GROUPS}
|
||||
saving={false}
|
||||
onCancel={() => {}}
|
||||
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%|不会遗漏/)
|
||||
})
|
||||
})
|
||||
@@ -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(
|
||||
<ToastProvider>
|
||||
<AutomationWorkspace dbReady={false} />
|
||||
</ToastProvider>
|
||||
)
|
||||
}
|
||||
|
||||
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/
|
||||
)
|
||||
})
|
||||
})
|
||||
@@ -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) => {
|
||||
|
||||
@@ -53,6 +53,14 @@ function fakeClock(step = 5): () => number {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 默认不真的等。
|
||||
*
|
||||
* 内置规则的 `replyDelaySeconds` 默认是 2 秒,如果每条例都真 sleep,
|
||||
* 这个文件会平白多出几十秒。**等待语义**由「回复前等待」那组用例显式覆盖。
|
||||
*/
|
||||
const noDelay = async (): Promise<void> => {}
|
||||
|
||||
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<WechatActionResult>
|
||||
generateReport?: () => Promise<AgentGroupReportResult>
|
||||
delay?: (ms: number) => Promise<void>
|
||||
}): {
|
||||
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()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
})
|
||||
@@ -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<void>
|
||||
/** 每次真正进入 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<AutomationRunResult> => {
|
||||
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<void>((resolve) => {
|
||||
release = resolve
|
||||
})
|
||||
let markStarted!: () => void
|
||||
const started = new Promise<void>((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<void>((resolve) => {
|
||||
release = resolve
|
||||
})
|
||||
let markStarted!: () => void
|
||||
const started = new Promise<void>((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
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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<string, Wcdb4Message[]>): {
|
||||
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)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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<WechatActionRequest> = {}
|
||||
): Promise<WechatActionResult> {
|
||||
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('自动化动作不允许执行')
|
||||
})
|
||||
})
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user