feat: 建立发送网关与动作审计

This commit is contained in:
Wxw-Gu
2026-09-02 18:00:39 +08:00
parent 76e9e683c4
commit e1f4ec79dd
10 changed files with 1411 additions and 88 deletions
+105 -19
View File
@@ -2,16 +2,17 @@ import { app, BrowserWindow } from 'electron'
import fs from 'fs-extra'
import path from 'path'
import * as chat from './chat-service'
import { personalWechatCapabilityService } from './personal-wechat-capability-service'
import { personalWechatSendService } from './personal-wechat-send-service'
import { wechatActionGateway, type WechatActionGateway } from './wechat-action-gateway'
import {
buildMemberLeftNotification,
findRemovedGroupMembers,
groupExitMemberName,
normalizeGroupExitNotificationTemplate,
renderGroupExitMonitorNotification,
validateGroupExitNotificationTemplate,
type GroupExitMonitorEvent,
type GroupExitMonitorMember,
type GroupExitNotificationStatus,
type GroupExitNotificationState,
type GroupExitMonitorState
} from '../../shared/group-exit-monitor'
@@ -39,7 +40,17 @@ const GROUP_READ_CONCURRENCY = 8
const DUPLICATE_WINDOW_MS = 2 * 60 * 1000
const MAX_EVENTS = 500
export interface GroupExitMonitorServiceDependencies {
actionGateway?: GroupExitActionGateway
}
type GroupExitActionGateway = Pick<WechatActionGateway, 'execute'> &
Partial<
Pick<WechatActionGateway, 'registerMemberEvent' | 'registerMemberEvents' | 'clearMemberEvents'>
>
class GroupExitMonitorService {
private readonly actionGateway: GroupExitActionGateway
private active = false
private nativeMonitorActive = false
private snapshots = new Map<string, GroupSnapshotRecord>()
@@ -59,6 +70,10 @@ class GroupExitMonitorService {
private scopeGeneration = 0
private notificationTemplate = normalizeGroupExitNotificationTemplate(undefined)
constructor(deps: GroupExitMonitorServiceDependencies = {}) {
this.actionGateway = deps.actionGateway || wechatActionGateway
}
getState(): GroupExitMonitorState {
this.ensureLoaded()
return {
@@ -84,6 +99,7 @@ class GroupExitMonitorService {
this.accountRoot &&
path.resolve(currentRoot) !== path.resolve(this.accountRoot)
) {
this.actionGateway.clearMemberEvents?.()
this.events = []
this.lastReadAt = 0
this.monitorSelectionConfigured = true
@@ -189,6 +205,7 @@ class GroupExitMonitorService {
clearEvents(): GroupExitMonitorState {
this.ensureLoaded()
this.events = []
this.actionGateway.clearMemberEvents?.()
this.lastReadAt = Date.now()
this.save()
this.broadcast()
@@ -225,6 +242,7 @@ class GroupExitMonitorService {
if (!currentRoomIds.has(roomId)) this.snapshots.delete(roomId)
}
const notificationTasks: Promise<void>[] = []
for (const next of groups) {
if (scopeGeneration !== this.scopeGeneration) return
if (next.membersValid === false) continue
@@ -246,10 +264,12 @@ class GroupExitMonitorService {
for (const member of removed) {
const event = this.recordExit(next, member, previous.members.length, next.members.length)
if (event && this.notificationRoomIds.has(next.roomId)) {
void this.notifyGroup(next, event)
notificationTasks.push(this.notifyGroup(next, event))
}
}
}
if (notificationTasks.length) await Promise.all(notificationTasks)
if (!this.active || scopeGeneration !== this.scopeGeneration) return
this.lastCheckedAt = Date.now()
this.save()
this.broadcast()
@@ -388,8 +408,10 @@ class GroupExitMonitorService {
currentCount,
delta: currentCount - previousCount,
message,
detectedAt
detectedAt,
notificationStatus: 'not_requested'
}
this.actionGateway.registerMemberEvent?.(event)
this.events = [event, ...this.events].slice(0, MAX_EVENTS)
console.log(
`[GroupMonitor] detected member exit roomId=${group.roomId} member=${member.wxid} ${previousCount}->${currentCount}`
@@ -401,29 +423,64 @@ class GroupExitMonitorService {
group: GroupSnapshotRecord,
event: GroupExitMonitorEvent
): Promise<void> {
event.notificationStatus = 'pending'
event.notification = { status: 'pending' }
this.save()
this.broadcast()
try {
const capability = await personalWechatCapabilityService.getPersonalWechatSendCapability()
if (!capability.ready || !capability.capabilities.text) {
console.warn(`[GroupMonitor] 跳过群聊通知 roomId=${group.roomId}: ${capability.message}`)
return
}
// 具体传输方式由发送服务按平台处理。
const result = await personalWechatSendService.send({
type: 'text',
to: group.roomId,
isGroup: true,
text: renderGroupExitMonitorNotification(event, this.notificationTemplate)
const result = await this.actionGateway.execute({
idempotencyKey: `member_left_notification:${event.id}`,
origin: 'member_monitor',
purpose: 'member_left_notification',
triggerType: 'automation',
sourceId: event.id,
recipient: {
type: 'group',
id: group.roomId,
name: group.groupName
},
content: {
type: 'text',
text: buildMemberLeftNotification(event, this.notificationTemplate)
},
metadata: {
memberId: event.memberWxid,
memberName: event.memberName,
detectedAt: event.detectedAt,
eventType: 'member_left',
eventRoomId: event.roomId
}
})
if (!result.success) {
const notification: GroupExitNotificationState = {
status: result.status,
actionId: result.actionId,
decision: result.decision,
...(result.errorCode ? { errorCode: result.errorCode } : {}),
...(result.reason ? { reason: result.reason } : {}),
startedAt: result.startedAt,
finishedAt: result.finishedAt
}
event.notificationStatus = result.status
event.notification = notification
if (result.status !== 'sent') {
console.warn(
`[GroupMonitor] 群聊通知发送失败 roomId=${group.roomId}: ${result.error || ''}`
`[GroupMonitor] 群聊通知未发送 roomId=${group.roomId} status=${result.status} code=${result.errorCode || ''}`
)
}
} catch (error) {
event.notificationStatus = 'failed'
event.notification = {
status: 'failed',
errorCode: 'UNKNOWN',
reason: error instanceof Error ? error.message : String(error)
}
console.warn(
`[GroupMonitor] 群聊通知异常 roomId=${group.roomId}:`,
error instanceof Error ? error.message : String(error)
)
} finally {
this.save()
this.broadcast()
}
}
@@ -437,6 +494,7 @@ class GroupExitMonitorService {
try {
const stored = fs.readJsonSync(this.filePath()) as StoredState
this.events = normalizeEvents(stored.events)
this.actionGateway.registerMemberEvents?.(this.events)
this.lastReadAt = Number(stored.lastReadAt) || 0
this.accountRoot = String(stored.accountRoot || '')
// 没有显式范围时按空范围处理,保留已有选择。
@@ -556,13 +614,41 @@ function normalizeEvents(
? Number(value.delta)
: currentCount - previousCount,
message: String(value.message || `${memberName}退出了${groupName}`),
detectedAt
detectedAt,
...(value.notificationStatus
? { notificationStatus: normalizeNotificationStatus(value.notificationStatus) }
: {}),
...(value.notification && typeof value.notification === 'object'
? { notification: normalizeNotification(value.notification) }
: {})
})
if (normalized.length >= MAX_EVENTS) break
}
return normalized
}
function normalizeNotificationStatus(value: unknown): GroupExitNotificationStatus {
const status = String(value || '').trim()
return status === 'pending' || status === 'sent' || status === 'blocked' || status === 'failed'
? status
: 'not_requested'
}
function normalizeNotification(value: object): GroupExitNotificationState {
const input = value as Partial<GroupExitNotificationState>
return {
status: normalizeNotificationStatus(input.status),
...(input.actionId ? { actionId: String(input.actionId) } : {}),
...(input.decision === 'allow' || input.decision === 'block'
? { decision: input.decision }
: {}),
...(input.errorCode ? { errorCode: String(input.errorCode) } : {}),
...(input.reason ? { reason: String(input.reason) } : {}),
...(input.startedAt ? { startedAt: String(input.startedAt) } : {}),
...(input.finishedAt ? { finishedAt: String(input.finishedAt) } : {})
}
}
function isContactEvent(rawPayload: string): boolean {
const payload = String(rawPayload || '').trim()
if (!payload) return false
+125 -69
View File
@@ -4,8 +4,10 @@ import { promises as fs } from 'fs'
import path from 'path'
import type {
PersonalWechatSendCapability,
PersonalWechatSendRequest
PersonalWechatSendRequest,
PersonalWechatSendResult
} from '../../shared/personal-wechat'
import type { WechatActionRequest, WechatActionResult } from '../../shared/wechat-action'
import type { AgentHubStatus } from '../../shared/agent-hub'
import type {
ScheduledReportCreateInput,
@@ -33,8 +35,8 @@ import {
import type { SaveGeneratedReportRequest } from '../../shared/report-history'
import { saveGeneratedReport } from '../report-history-service'
import { generateAgentGroupReport } from './agent-group-report-service'
import { personalWechatSendService } from './personal-wechat-send-service'
import { personalWechatCapabilityService } from './personal-wechat-capability-service'
import { WechatActionGateway, wechatActionGateway } from './wechat-action-gateway'
import { getContactAvatars, isReady as isChatReady, resolveMd5 } from './chat-service'
import { agentHubService, type AgentHubNotificationResult } from './agent-hub-service'
@@ -55,7 +57,8 @@ export interface ScheduledReportDependencies {
saveGeneratedReport: (
request: SaveGeneratedReportRequest
) => ReturnType<typeof saveGeneratedReport>
send: (request: PersonalWechatSendRequest) => ReturnType<typeof personalWechatSendService.send>
send: (request: PersonalWechatSendRequest) => Promise<PersonalWechatSendResult>
sendAction?: (input: ScheduledReportSendActionInput) => Promise<WechatActionResult>
sendNotification: (input: { to?: string; text: string }) => Promise<AgentHubNotificationResult>
getNotificationRecipient: () => string | undefined
getAgentHubStatus: () => AgentHubStatus
@@ -65,11 +68,43 @@ export interface ScheduledReportDependencies {
now?: () => Date
}
export interface ScheduledReportSendActionInput {
target: string
filePath: string
executionId: string
triggerType: 'scheduled' | 'manual'
retryCount?: number
taskId?: string
}
function buildScheduledReportActionRequest(
input: ScheduledReportSendActionInput
): WechatActionRequest {
return {
idempotencyKey:
input.retryCount && input.retryCount > 0
? `scheduled_report:${input.executionId}:retry:${input.retryCount}`
: `scheduled_report:${input.executionId}`,
origin: 'scheduled_report',
purpose: 'scheduled_report',
triggerType: input.triggerType === 'scheduled' ? 'automation' : 'user',
sourceId: input.executionId,
executionId: input.executionId,
recipient: { type: 'group', id: input.target },
content: { type: 'image', path: input.filePath },
metadata: {
taskId: input.taskId,
retryCount: input.retryCount || 0
}
}
}
const defaultDependencies = (): ScheduledReportDependencies => ({
getCapability: () => personalWechatCapabilityService.getPersonalWechatSendCapability(),
generateReport: generateAgentGroupReport,
saveGeneratedReport,
send: (request) => personalWechatSendService.send(request),
send: (request) => sendRequestThroughGateway(request),
sendAction: (input) => wechatActionGateway.execute(buildScheduledReportActionRequest(input)),
sendNotification: (input) => agentHubService.sendNotification(input),
getNotificationRecipient: () => agentHubService.getNotificationRecipient(),
getAgentHubStatus: () => agentHubService.getStatus(),
@@ -78,6 +113,35 @@ const defaultDependencies = (): ScheduledReportDependencies => ({
isDatabaseReady: () => isChatReady()
})
async function sendRequestThroughGateway(
request: PersonalWechatSendRequest
): Promise<PersonalWechatSendResult> {
const content =
request.type === 'text'
? { type: 'text' as const, text: request.text }
: { type: request.type, path: request.filePath }
const action = await wechatActionGateway.execute({
origin: 'scheduled_report',
purpose: 'scheduled_report',
triggerType: 'automation',
recipient: {
type: request.isGroup ? 'group' : 'contact',
id: request.to
},
content
})
const sendResult = action.sendResult
const status =
sendResult && typeof sendResult === 'object' && 'status' in sendResult
? (sendResult as { status: PersonalWechatSendResult['status'] }).status
: ({} as PersonalWechatSendResult['status'])
return {
success: action.status === 'sent',
status,
...(action.reason || action.errorCode ? { error: action.reason || action.errorCode } : {})
}
}
const rangeValues = new Set<ScheduledReportRange>(['today', 'yesterday', '7days', 'recent24h'])
const reportRangeLabel = (range: ScheduledReportRange): string =>
@@ -148,7 +212,21 @@ export class ScheduledReportService {
private readonly retrying = new Map<string, Promise<ScheduledReportExecution>>()
constructor(deps?: Partial<ScheduledReportDependencies>) {
this.deps = { ...defaultDependencies(), ...deps }
const defaults = defaultDependencies()
const merged = { ...defaults, ...deps }
if (deps?.send && !deps.sendAction) {
// 保留注入的传输实现,方便测试和内嵌场景,同时让日报发送继续使用统一的发送入口。
const legacyGateway = new WechatActionGateway({
getCapability: merged.getCapability,
send: merged.send,
getUserDataPath: () => merged.storageDir,
...(merged.now ? { now: merged.now } : {})
})
merged.sendAction = (input) => legacyGateway.execute(buildScheduledReportActionRequest(input))
} else if (!merged.sendAction) {
merged.sendAction = defaults.sendAction
}
this.deps = merged
}
async start(): Promise<void> {
@@ -626,8 +704,6 @@ export class ScheduledReportService {
try {
await update({ currentStage: 'precheck' })
let capability: PersonalWechatSendCapability | null = null
let capabilityCheckError = ''
await update({ currentStage: 'data' })
await update({ currentStage: 'ai' })
@@ -717,27 +793,6 @@ export class ScheduledReportService {
}
await update({ reportId, htmlPath, pngPath, currentStage: 'send' })
try {
capability = await this.deps.getCapability()
} catch (error) {
capabilityCheckError = error instanceof Error ? error.message : String(error)
}
if (!capability?.ready || !capability.capabilities.image) {
const technicalMessage =
capabilityCheckError ||
capability?.error ||
capability?.message ||
'个人微信发送能力不可用'
return finishError(
technicalMessage,
'send',
'waiting_to_send',
{ sendStatus: 'unavailable', sendError: technicalMessage },
'partial_success',
'WECHAT_SEND_UNAVAILABLE'
)
}
const target = this.resolveTarget(task)
if (!target) {
return finishError(
@@ -750,13 +805,14 @@ export class ScheduledReportService {
)
}
await update({ sendTarget: target, currentStage: 'send' })
let sent: Awaited<ReturnType<ScheduledReportDependencies['send']>>
let action: WechatActionResult
try {
sent = await this.deps.send({
type: 'image',
to: target,
isGroup: true,
filePath: pngPath
action = await this.deps.sendAction!({
target,
filePath: pngPath,
executionId: execution.id,
triggerType,
taskId: task.id
})
} catch (error) {
return finishError(
@@ -771,14 +827,21 @@ export class ScheduledReportService {
'WECHAT_SEND_FAILED'
)
}
if (!sent.success) {
if (action.status !== 'sent') {
const unavailable =
action.errorCode === 'SEND_CAPABILITY_UNAVAILABLE' ||
action.errorCode === 'SEND_NOT_READY'
const actionError = action.reason || action.errorCode || '微信发送失败'
return finishError(
sent.error || '微信发送失败',
actionError,
'send',
unavailable ? 'waiting_to_send' : 'partial_success',
{
sendStatus: unavailable ? 'unavailable' : 'failed',
sendError: actionError
},
'partial_success',
{ sendStatus: 'failed', sendError: sent.error || '微信发送失败' },
'partial_success',
'WECHAT_SEND_FAILED'
unavailable ? 'WECHAT_SEND_UNAVAILABLE' : 'WECHAT_SEND_FAILED'
)
}
return this.finalizeExecution(task, execution, {
@@ -844,38 +907,28 @@ export class ScheduledReportService {
)
}
let capability: PersonalWechatSendCapability | null = null
try {
capability = await this.deps.getCapability()
} catch (error) {
return finishRetryError(error, 'WECHAT_SEND_UNAVAILABLE', 'unavailable', 'waiting_to_send')
}
if (!capability.ready || !capability.capabilities.image) {
return finishRetryError(
capability.error || capability.message || '个人微信发送能力不可用',
'WECHAT_SEND_UNAVAILABLE',
'unavailable',
'waiting_to_send'
)
}
const target = working.sendTarget || this.resolveTarget(task)
if (!target) {
return finishRetryError('未找到指定微信群', 'WECHAT_SEND_FAILED', 'failed', 'partial_success')
}
try {
const sent = await this.deps.send({
type: 'image',
to: target,
isGroup: true,
filePath: working.pngPath
const action = await this.deps.sendAction!({
target,
filePath: working.pngPath,
executionId: working.id,
triggerType: working.triggerType || 'scheduled',
retryCount,
taskId: task.id
})
if (!sent.success) {
if (action.status !== 'sent') {
const unavailable =
action.errorCode === 'SEND_CAPABILITY_UNAVAILABLE' ||
action.errorCode === 'SEND_NOT_READY'
return finishRetryError(
sent.error || '微信发送失败',
'WECHAT_SEND_FAILED',
'failed',
'partial_success'
action.reason || action.errorCode || '微信发送失败',
unavailable ? 'WECHAT_SEND_UNAVAILABLE' : 'WECHAT_SEND_FAILED',
unavailable ? 'unavailable' : 'failed',
unavailable ? 'waiting_to_send' : 'partial_success'
)
}
} catch (error) {
@@ -932,18 +985,21 @@ export class ScheduledReportService {
...patch,
finishedAt: (this.deps.now?.() || new Date()).toISOString()
}
await this.persistExecution(completed)
// 让执行记录在通知完成前保持进行中,避免日报已结束但通知状态尚未更新。
const notified = notification
? await this.notifyExecution(task, completed, notification)
: completed
await this.persistExecution(notified)
const taskIndex = this.tasks!.findIndex((item) => item.id === task.id)
if (taskIndex >= 0) {
this.tasks![taskIndex] = {
...this.tasks![taskIndex],
lastRunAt: completed.finishedAt,
updatedAt: completed.finishedAt!
lastRunAt: notified.finishedAt,
updatedAt: notified.finishedAt!
}
await this.saveTasks()
}
if (notification) return this.notifyExecution(task, completed, notification)
return { ...completed }
return { ...notified }
}
private async persistExecution(execution: ScheduledReportExecution): Promise<void> {
+619
View File
@@ -0,0 +1,619 @@
import { app } from 'electron'
import { createHash, randomUUID } from 'crypto'
import fs from 'fs-extra'
import path from 'path'
import type {
PersonalWechatSendCapability,
PersonalWechatSendRequest,
PersonalWechatSendResult
} from '../../shared/personal-wechat'
import type {
PolicyDecision,
WechatActionAuditRecord,
WechatActionContent,
WechatActionErrorCode,
WechatActionMemberEventReference,
WechatActionRequest,
WechatActionResult
} from '../../shared/wechat-action'
import { personalWechatCapabilityService } from './personal-wechat-capability-service'
import { personalWechatSendService } from './personal-wechat-send-service'
const MAX_AUDIT_RECORDS = 500
const MAX_CONTENT_PREVIEW_LENGTH = 240
const AUTOMATION_PURPOSE_ALLOWLIST = new Set(['scheduled_report', 'member_left_notification'])
export interface WechatActionGatewayDependencies {
getCapability?: () => Promise<PersonalWechatSendCapability>
send?: (request: PersonalWechatSendRequest) => Promise<PersonalWechatSendResult>
getMemberEvent?: (
sourceId: string
) =>
| WechatActionMemberEventReference
| Promise<WechatActionMemberEventReference | undefined>
| undefined
checkEntitlement?: (request: WechatActionRequest) => boolean | Promise<boolean>
getUserDataPath?: () => string
now?: () => Date
}
interface LoadedAuditState {
path: string
records: WechatActionAuditRecord[]
}
export interface WechatActionPolicyContext {
memberEvent?: WechatActionMemberEventReference
}
const defaultDependencies = (): Required<
Pick<WechatActionGatewayDependencies, 'getCapability' | 'send' | 'getUserDataPath' | 'now'>
> => ({
getCapability: () => personalWechatCapabilityService.getPersonalWechatSendCapability(),
send: (request) => personalWechatSendService.send(request),
getUserDataPath: () => app.getPath('userData'),
now: () => new Date()
})
/**
* 统一管理微信发送操作,在真正发送前完成必要检查;
* 具体发送由底层服务处理,使用者只需提供发送内容和对象。
*/
export class WechatActionGateway {
private readonly deps: Required<
Pick<WechatActionGatewayDependencies, 'getCapability' | 'send' | 'getUserDataPath' | 'now'>
> &
Omit<WechatActionGatewayDependencies, 'getCapability' | 'send' | 'getUserDataPath' | 'now'>
private readonly memberEvents = new Map<string, WechatActionMemberEventReference>()
private readonly inFlight = new Map<string, Promise<WechatActionResult>>()
private auditState: LoadedAuditState | null = null
constructor(deps: WechatActionGatewayDependencies = {}) {
this.deps = { ...defaultDependencies(), ...deps }
}
/** 记录退群事件,发送通知前可确认事件所属群聊。 */
registerMemberEvent(event: WechatActionMemberEventReference): void {
const id = String(event?.id || '').trim()
const roomId = String(event?.roomId || '').trim()
if (!id || !roomId) return
this.memberEvents.set(id, { id, roomId })
}
registerMemberEvents(events: WechatActionMemberEventReference[]): void {
for (const event of events) this.registerMemberEvent(event)
}
clearMemberEvents(): void {
this.memberEvents.clear()
}
async execute(request: WechatActionRequest): Promise<WechatActionResult> {
const normalized = normalizeActionRequest(request)
const idempotencyKey = normalized ? this.idempotencyKey(normalized) : undefined
if (idempotencyKey) {
const existing = await this.findExisting(idempotencyKey)
if (existing) return existing
const running = this.inFlight.get(idempotencyKey)
if (running) return running
}
const promise = this.executeOnce(request, normalized, idempotencyKey)
if (idempotencyKey) this.inFlight.set(idempotencyKey, promise)
try {
return await promise
} finally {
if (idempotencyKey && this.inFlight.get(idempotencyKey) === promise) {
this.inFlight.delete(idempotencyKey)
}
}
}
private async executeOnce(
originalRequest: WechatActionRequest,
request: WechatActionRequest | null,
idempotencyKey?: string
): Promise<WechatActionResult> {
const startedAt = this.deps.now().toISOString()
const actionId = String(originalRequest?.id || '').trim() || randomUUID()
if (!request) {
return this.finishBlocked(
actionId,
originalRequest,
idempotencyKey,
startedAt,
validationErrorCode(originalRequest),
validationErrorReason(originalRequest)
)
}
const eventContext = await this.resolveEventContext(request)
const policy = evaluateWechatActionPolicy(request, eventContext)
if (policy.decision !== 'allow') {
return this.finishBlocked(
actionId,
request,
idempotencyKey,
startedAt,
policy.reasonCode || 'POLICY_BLOCKED',
policy.reason || '该发送动作未通过策略检查'
)
}
let entitled = true
try {
entitled = this.deps.checkEntitlement
? await this.deps.checkEntitlement(request)
: await checkEntitlement(request)
} catch (error) {
entitled = false
return this.finishBlocked(
actionId,
request,
idempotencyKey,
startedAt,
'POLICY_BLOCKED',
error instanceof Error ? error.message : String(error)
)
}
if (!entitled) {
return this.finishBlocked(
actionId,
request,
idempotencyKey,
startedAt,
'POLICY_BLOCKED',
'该发送动作未获得当前账号的使用权限'
)
}
const capabilityType = request.content.type
let capability: PersonalWechatSendCapability | undefined
try {
capability = await this.deps.getCapability()
} catch (error) {
return this.finishFailed(
actionId,
request,
idempotencyKey,
startedAt,
'SEND_CAPABILITY_UNAVAILABLE',
error instanceof Error ? error.message : String(error)
)
}
if (!capability || !capability.ready) {
return this.finishFailed(
actionId,
request,
idempotencyKey,
startedAt,
'SEND_CAPABILITY_UNAVAILABLE',
capability?.error || capability?.message || '个人微信发送能力不可用'
)
}
if (!capability.capabilities?.[capabilityType]) {
return this.finishFailed(
actionId,
request,
idempotencyKey,
startedAt,
'SEND_NOT_READY',
capability.message || `个人微信尚未准备好发送${contentTypeLabel(capabilityType)}`
)
}
let sendResult: PersonalWechatSendResult
try {
sendResult = await this.deps.send(toPersonalWechatSendRequest(request))
} catch (error) {
return this.finishFailed(
actionId,
request,
idempotencyKey,
startedAt,
'SEND_FAILED',
error instanceof Error ? error.message : String(error)
)
}
if (!sendResult.success) {
return this.finishFailed(
actionId,
request,
idempotencyKey,
startedAt,
'SEND_FAILED',
sendResult.error || '微信发送失败',
sendResult
)
}
return this.finishSent(actionId, request, idempotencyKey, startedAt, sendResult)
}
private async resolveEventContext(
request: WechatActionRequest
): Promise<WechatActionPolicyContext> {
if (request.purpose !== 'member_left_notification') return {}
const sourceId = String(request.sourceId || '').trim()
if (!sourceId) return {}
let memberEvent: WechatActionMemberEventReference | undefined = this.memberEvents.get(sourceId)
if (!memberEvent && this.deps.getMemberEvent) {
let resolved: WechatActionMemberEventReference | undefined
try {
resolved = await this.deps.getMemberEvent(sourceId)
} catch {
resolved = undefined
}
if (resolved && String(resolved.id || '').trim() === sourceId) {
memberEvent = { id: String(resolved.id), roomId: String(resolved.roomId) }
this.registerMemberEvent(memberEvent)
}
}
return memberEvent ? { memberEvent } : {}
}
private async findExisting(idempotencyKey: string): Promise<WechatActionResult | undefined> {
const state = this.loadAuditState()
const existing = state.records.find((record) => record.idempotencyKey === idempotencyKey)
if (!existing) return undefined
return {
actionId: existing.actionId,
status: existing.sendStatus,
decision: existing.decision,
...(existing.errorCode ? { errorCode: existing.errorCode } : {}),
...(existing.decisionReason ? { reason: existing.decisionReason } : {}),
startedAt: existing.startedAt,
finishedAt: existing.finishedAt
}
}
private idempotencyKey(request: WechatActionRequest): string | undefined {
const explicit = String(request.idempotencyKey || '').trim()
if (explicit) return explicit
if (request.purpose === 'member_left_notification' && request.sourceId) {
return `member_left_notification:${request.sourceId}`
}
if (request.origin === 'scheduled_report' && request.executionId) {
return `scheduled_report:${request.executionId}`
}
return undefined
}
private finishSent(
actionId: string,
request: WechatActionRequest,
idempotencyKey: string | undefined,
startedAt: string,
sendResult: PersonalWechatSendResult
): WechatActionResult {
const finishedAt = this.deps.now().toISOString()
const result: WechatActionResult = {
actionId,
status: 'sent',
decision: 'allow',
startedAt,
finishedAt,
sendResult
}
this.writeAudit(request, idempotencyKey, result)
return result
}
private finishFailed(
actionId: string,
request: WechatActionRequest,
idempotencyKey: string | undefined,
startedAt: string,
errorCode: WechatActionErrorCode,
reason: string,
sendResult?: PersonalWechatSendResult
): WechatActionResult {
const finishedAt = this.deps.now().toISOString()
const result: WechatActionResult = {
actionId,
status: 'failed',
decision: 'allow',
errorCode,
reason,
startedAt,
finishedAt,
...(sendResult ? { sendResult } : {})
}
this.writeAudit(request, idempotencyKey, result)
return result
}
private finishBlocked(
actionId: string,
request: WechatActionRequest,
idempotencyKey: string | undefined,
startedAt: string,
errorCode: WechatActionErrorCode,
reason: string
): WechatActionResult {
const finishedAt = this.deps.now().toISOString()
const result: WechatActionResult = {
actionId,
status: 'blocked',
decision: 'block',
errorCode,
reason,
startedAt,
finishedAt
}
this.writeAudit(request, idempotencyKey, result)
return result
}
private loadAuditState(): LoadedAuditState {
const filePath = this.auditFilePath()
if (this.auditState?.path === filePath) return this.auditState
let records: WechatActionAuditRecord[] = []
try {
const value = fs.readJsonSync(filePath) as unknown
if (Array.isArray(value)) records = normalizeAuditRecords(value)
else if (
value &&
typeof value === 'object' &&
Array.isArray((value as { actions?: unknown }).actions)
) {
records = normalizeAuditRecords((value as { actions: unknown[] }).actions)
}
} catch {
records = []
}
this.auditState = { path: filePath, records }
return this.auditState
}
private writeAudit(
request: WechatActionRequest,
idempotencyKey: string | undefined,
result: WechatActionResult
): void {
if (
!request ||
!request.recipient ||
(request.recipient.type !== 'group' && request.recipient.type !== 'contact') ||
!request.content ||
(request.content.type !== 'text' &&
request.content.type !== 'image' &&
request.content.type !== 'voice')
) {
return
}
const state = this.loadAuditState()
const record: WechatActionAuditRecord = {
actionId: result.actionId,
...(idempotencyKey ? { idempotencyKey } : {}),
origin: request.origin,
purpose: request.purpose,
triggerType: request.triggerType,
...(request.sourceId ? { sourceId: request.sourceId } : {}),
...(request.executionId ? { executionId: request.executionId } : {}),
recipientType: request.recipient.type,
recipientId: request.recipient.id,
contentType: request.content.type,
...contentAudit(request.content),
createdAt: result.startedAt,
startedAt: result.startedAt,
finishedAt: result.finishedAt,
decision: result.decision,
...(result.reason ? { decisionReason: result.reason } : {}),
sendStatus: result.status,
...(result.errorCode ? { errorCode: result.errorCode } : {})
}
const withoutKey = state.records.filter((item) => item.actionId !== record.actionId)
state.records = [record, ...withoutKey].slice(0, MAX_AUDIT_RECORDS)
try {
fs.ensureDirSync(path.dirname(state.path))
fs.writeJsonSync(state.path, state.records, { spaces: 2 })
} catch (error) {
console.warn('[WechatActionGateway] 保存 Action 审计失败:', error)
}
}
private auditFilePath(): string {
return path.join(this.deps.getUserDataPath(), 'actions', 'wechat-actions.json')
}
}
export function evaluateWechatActionPolicy(
request: WechatActionRequest,
context: WechatActionPolicyContext = {}
): PolicyDecision {
if (!request || typeof request !== 'object') {
return {
decision: 'block',
source: 'deterministic',
reasonCode: 'INVALID_REQUEST',
reason: '发送动作请求格式无效'
}
}
const recipient = request.recipient
if (
!recipient ||
typeof recipient !== 'object' ||
(recipient.type !== 'group' && recipient.type !== 'contact') ||
!String(recipient.id || '').trim()
) {
return {
decision: 'block',
source: 'deterministic',
reasonCode: 'INVALID_RECIPIENT',
reason: '发送动作必须指定有效的收件人'
}
}
if (request.triggerType === 'automation' && !AUTOMATION_PURPOSE_ALLOWLIST.has(request.purpose)) {
return {
decision: 'block',
source: 'deterministic',
reasonCode: 'ACTION_NOT_ALLOWED',
reason: `自动化动作不允许执行 purpose=${request.purpose}`
}
}
if (request.purpose === 'member_left_notification') {
if (
!request.sourceId ||
!context.memberEvent ||
context.memberEvent.id !== request.sourceId ||
!context.memberEvent.roomId
) {
return {
decision: 'block',
source: 'deterministic',
reasonCode: 'INVALID_REQUEST',
reason: '退群通知必须关联已记录的退群事件'
}
}
if (request.recipient.type !== 'group' || request.recipient.id !== context.memberEvent.roomId) {
return {
decision: 'block',
source: 'deterministic',
reasonCode: 'RECIPIENT_SCOPE_VIOLATION',
reason: '退群通知只能发送回原事件所在群聊'
}
}
}
return { decision: 'allow', source: 'deterministic' }
}
/** 发送前会根据固定规则检查是否允许发送。 */
export function shouldUseAiPolicy(request: WechatActionRequest): boolean {
void request
return false
}
/** 当前默认允许发送,后续可在这里补充权限限制。 */
export function checkEntitlement(request: WechatActionRequest): boolean {
void request
return true
}
export function normalizeActionRequest(value: unknown): WechatActionRequest | null {
if (!value || typeof value !== 'object') return null
const input = value as Partial<WechatActionRequest>
const origin = String(input.origin || '').trim()
const purpose = String(input.purpose || '').trim()
const triggerType = input.triggerType
const recipient = input.recipient
const content = input.content
if (!origin || !purpose || (triggerType !== 'automation' && triggerType !== 'user')) return null
if (!recipient || typeof recipient !== 'object') return null
const recipientType = (recipient as { type?: unknown }).type
const recipientId = String((recipient as { id?: unknown }).id || '').trim()
if (recipientType !== 'group' && recipientType !== 'contact') return null
if (!recipientId) return null
if (!content || typeof content !== 'object') return null
const contentType = (content as { type?: unknown }).type
if (contentType === 'text') {
const text = String((content as { text?: unknown }).text || '').trim()
if (!text || text.length > 2_000) return null
return {
...input,
origin,
purpose,
triggerType,
...(input.id ? { id: String(input.id).trim() } : {}),
...(input.idempotencyKey ? { idempotencyKey: String(input.idempotencyKey).trim() } : {}),
...(input.sourceId ? { sourceId: String(input.sourceId).trim() } : {}),
...(input.executionId ? { executionId: String(input.executionId).trim() } : {}),
recipient: {
type: recipientType,
id: recipientId,
...((recipient as { name?: unknown }).name
? { name: String((recipient as { name?: unknown }).name).trim() }
: {})
},
content: { type: 'text', text }
} as WechatActionRequest
}
if (contentType !== 'image' && contentType !== 'voice') return null
const filePath = String((content as { path?: unknown }).path || '').trim()
if (!filePath) return null
return {
...input,
origin,
purpose,
triggerType,
...(input.id ? { id: String(input.id).trim() } : {}),
...(input.idempotencyKey ? { idempotencyKey: String(input.idempotencyKey).trim() } : {}),
...(input.sourceId ? { sourceId: String(input.sourceId).trim() } : {}),
...(input.executionId ? { executionId: String(input.executionId).trim() } : {}),
recipient: {
type: recipientType,
id: recipientId,
...((recipient as { name?: unknown }).name
? { name: String((recipient as { name?: unknown }).name).trim() }
: {})
},
content: { type: contentType, path: filePath }
} as WechatActionRequest
}
function validationErrorCode(value: unknown): WechatActionErrorCode {
if (!value || typeof value !== 'object') return 'INVALID_REQUEST'
const recipient = (value as { recipient?: unknown }).recipient
if (!recipient || typeof recipient !== 'object') return 'INVALID_RECIPIENT'
const recipientValue = recipient as { type?: unknown; id?: unknown }
if (
(recipientValue.type !== 'group' && recipientValue.type !== 'contact') ||
!String(recipientValue.id || '').trim()
) {
return 'INVALID_RECIPIENT'
}
return 'INVALID_REQUEST'
}
function validationErrorReason(value: unknown): string {
return validationErrorCode(value) === 'INVALID_RECIPIENT'
? '发送动作必须指定有效的收件人'
: '发送动作请求格式无效'
}
function toPersonalWechatSendRequest(request: WechatActionRequest): PersonalWechatSendRequest {
const base = {
to: request.recipient.id,
isGroup: request.recipient.type === 'group'
}
if (request.content.type === 'text') return { ...base, type: 'text', text: request.content.text }
if (request.content.type === 'image') {
return { ...base, type: 'image', filePath: request.content.path }
}
return { ...base, type: 'voice', filePath: request.content.path }
}
function contentAudit(
content: WechatActionContent
): Pick<WechatActionAuditRecord, 'contentPreview' | 'contentHash'> {
const raw = content.type === 'text' ? content.text : content.path
const preview = content.type === 'text' ? raw : path.basename(raw)
return {
contentPreview: preview.slice(0, MAX_CONTENT_PREVIEW_LENGTH),
contentHash: createHash('sha256').update(raw).digest('hex')
}
}
function contentTypeLabel(type: WechatActionContent['type']): string {
return type === 'text' ? '文字' : type === 'image' ? '图片' : '语音'
}
function normalizeAuditRecords(values: unknown[]): WechatActionAuditRecord[] {
return values
.filter((value): value is WechatActionAuditRecord => {
if (!value || typeof value !== 'object') return false
const record = value as Partial<WechatActionAuditRecord>
return Boolean(
String(record.actionId || '').trim() &&
String(record.recipientId || '').trim() &&
String(record.contentType || '').trim() &&
String(record.startedAt || '').trim() &&
String(record.finishedAt || '').trim()
)
})
.slice(0, MAX_AUDIT_RECORDS)
}
export const wechatActionGateway = new WechatActionGateway()
/** 兼容性名称,方便已有调用继续使用。 */
export const personalWechatActionService = wechatActionGateway
@@ -80,6 +80,38 @@ const sendCapabilityLabel = (capability: PersonalWechatSendCapability | null): s
const sendCapabilityTone = (capability: PersonalWechatSendCapability | null): string =>
capability?.ready && capability.capabilities.text ? 'ready' : 'unready'
function EventNotificationStatus({
event
}: {
event: GroupExitMonitorState['events'][number]
}): React.ReactElement {
const status = event.notificationStatus || event.notification?.status || 'not_requested'
const details =
status === 'sent'
? { label: '已通知当前群聊', tone: 'sent', description: '' }
: status === 'pending'
? { label: '正在通知当前群聊', tone: 'pending', description: '' }
: status === 'blocked'
? { label: '通知已拦截', tone: 'blocked', description: '该操作未通过发送策略检查。' }
: status === 'failed'
? {
label: '通知未发送',
tone: 'failed',
description:
event.notification?.errorCode === 'SEND_CAPABILITY_UNAVAILABLE' ||
event.notification?.errorCode === 'SEND_NOT_READY'
? '当前微信发送能力不可用,退群事件已正常记录。'
: '退群事件已正常记录,但通知发送失败。'
}
: { label: '仅记录', tone: 'recorded', description: '' }
return (
<div className={`exit-monitor-event-notification ${details.tone}`}>
<span>{details.label}</span>
{details.description ? <small>{details.description}</small> : null}
</div>
)
}
function SendCapabilityStatus({
capability,
className = ''
@@ -783,6 +815,7 @@ export function GroupExitMonitorWorkspace({
群人数 {event.previousCount} 人 → {event.currentCount} 人
</p>
) : null}
<EventNotificationStatus event={event} />
<dl className="exit-monitor-event-details">
<div>
<dt>微信名</dt>
@@ -350,6 +350,40 @@
color: var(--wxex-danger);
}
.exit-monitor-event-notification {
display: flex;
flex-wrap: wrap;
align-items: baseline;
gap: 6px;
margin-top: 7px;
color: var(--wxex-text-secondary);
font-size: 12px;
line-height: 18px;
}
.exit-monitor-event-notification span {
font-weight: 600;
}
.exit-monitor-event-notification small {
color: var(--wxex-text-muted);
font-size: 11px;
line-height: 18px;
}
.exit-monitor-event-notification.sent {
color: var(--wxex-success);
}
.exit-monitor-event-notification.blocked,
.exit-monitor-event-notification.failed {
color: var(--wxex-danger);
}
.exit-monitor-event-notification.pending {
color: var(--wxex-brand);
}
.exit-monitor-event-details {
display: grid;
grid-template-columns: repeat(2, minmax(0, 1fr));
+30
View File
@@ -1,3 +1,5 @@
import type { WechatActionResult } from './wechat-action'
export interface GroupExitMonitorMember {
wxid: string
nickname?: string
@@ -7,6 +9,23 @@ export interface GroupExitMonitorMember {
avatar?: string
}
export type GroupExitNotificationStatus =
| 'not_requested'
| 'pending'
| 'sent'
| 'blocked'
| 'failed'
export interface GroupExitNotificationState {
status: GroupExitNotificationStatus
actionId?: string
decision?: WechatActionResult['decision']
errorCode?: string
reason?: string
startedAt?: string
finishedAt?: string
}
export interface GroupExitMonitorEvent {
id: string
contactId: string
@@ -25,6 +44,9 @@ export interface GroupExitMonitorEvent {
delta: number
message: string
detectedAt: number
/** 可选通知的结果;退群检测与通知发送彼此独立。 */
notificationStatus?: GroupExitNotificationStatus
notification?: GroupExitNotificationState
}
const LEGACY_GROUP_EXIT_NOTIFICATION_TEMPLATE = [
@@ -142,6 +164,14 @@ export function renderGroupExitMonitorNotification(
)
}
/** 主进程中的辅助函数,供生成通知和预览时共用。 */
export function buildMemberLeftNotification(
event: GroupExitMonitorEvent,
template = GROUP_EXIT_NOTIFICATION_TEMPLATE
): string {
return renderGroupExitMonitorNotification(event, template)
}
export interface GroupExitMonitorState {
events: GroupExitMonitorEvent[]
running: boolean
+100
View File
@@ -0,0 +1,100 @@
export type WechatActionOrigin =
| 'member_monitor'
| 'scheduled_report'
| 'user_tts'
| 'unknown'
| (string & {})
export type WechatActionPurpose =
| 'member_left_notification'
| 'scheduled_report'
| 'tts_voice'
| (string & {})
export type WechatActionTriggerType = 'automation' | 'user'
export interface WechatActionRecipient {
type: 'group' | 'contact'
id: string
name?: string
}
export type WechatActionContent =
| { type: 'text'; text: string }
| { type: 'image'; path: string }
| { type: 'voice'; path: string }
export interface WechatActionRequest {
/** 可选的操作编号;未提供时会自动生成。 */
id?: string
/** 用于避免同一自动操作被重复执行的标识。 */
idempotencyKey?: string
origin: WechatActionOrigin
purpose: WechatActionPurpose
triggerType: WechatActionTriggerType
sourceId?: string
executionId?: string
recipient: WechatActionRecipient
content: WechatActionContent
metadata?: Record<string, unknown>
}
export type WechatActionStatus = 'sent' | 'blocked' | 'failed'
export type WechatActionDecision = 'allow' | 'block'
export type WechatActionErrorCode =
| 'INVALID_REQUEST'
| 'INVALID_RECIPIENT'
| 'ACTION_NOT_ALLOWED'
| 'RECIPIENT_SCOPE_VIOLATION'
| 'SEND_CAPABILITY_UNAVAILABLE'
| 'SEND_NOT_READY'
| 'SEND_FAILED'
| 'POLICY_BLOCKED'
| 'UNKNOWN'
| (string & {})
export interface WechatActionResult {
actionId: string
status: WechatActionStatus
decision: WechatActionDecision
errorCode?: WechatActionErrorCode
reason?: string
startedAt: string
finishedAt: string
sendResult?: unknown
}
export interface PolicyDecision {
decision: 'allow' | 'block' | 'require_review'
source: 'deterministic' | 'ai'
reasonCode?: WechatActionErrorCode
reason?: string
}
export interface WechatActionMemberEventReference {
id: string
roomId: string
}
export interface WechatActionAuditRecord {
actionId: string
idempotencyKey?: string
origin: WechatActionOrigin
purpose: WechatActionPurpose
triggerType: WechatActionTriggerType
sourceId?: string
executionId?: string
recipientType: WechatActionRecipient['type']
recipientId: string
contentType: WechatActionContent['type']
contentPreview?: string
contentHash?: string
createdAt: string
startedAt: string
finishedAt: string
decision: WechatActionDecision
decisionReason?: string
sendStatus: WechatActionStatus
errorCode?: WechatActionErrorCode
}
@@ -159,4 +159,84 @@ describe('GroupExitMonitorService', () => {
text: '退群: 另一个人/未设置/wxid_other'
})
})
it('records an exit without creating a send action when notification is disabled', async () => {
installGroupDb()
mocks.chat.getGroupSnapshotAsync
.mockResolvedValueOnce({
roomId: 'room@chatroom',
members: [member, { wxid: 'wxid_other', wechatNickname: '另一个人' }]
})
.mockResolvedValueOnce({ roomId: 'room@chatroom', members: [member] })
const service = new GroupExitMonitorService()
service.setMonitoredRoomIds(['room@chatroom'])
await service.start(true)
const state = await service.checkNow()
expect(state.events[0]).toMatchObject({
memberWxid: 'wxid_other',
notificationStatus: 'not_requested'
})
expect(mocks.sender.send).not.toHaveBeenCalled()
expect(mocks.capability.getPersonalWechatSendCapability).not.toHaveBeenCalled()
})
it('keeps the member event when the Gateway reports unavailable capability', async () => {
installGroupDb()
mocks.capability.getPersonalWechatSendCapability.mockResolvedValue({
ready: false,
supported: true,
status: 'needs_binding',
capabilities: { text: false, image: false, voice: false },
senderStatus: {},
message: '当前微信发送能力不可用'
})
mocks.chat.getGroupSnapshotAsync
.mockResolvedValueOnce({
roomId: 'room@chatroom',
members: [member, { wxid: 'wxid_other', wechatNickname: '另一个人' }]
})
.mockResolvedValueOnce({ roomId: 'room@chatroom', members: [member] })
const service = new GroupExitMonitorService()
service.setMonitoredRoomIds(['room@chatroom'], ['room@chatroom'])
await service.start(true)
const state = await service.checkNow()
expect(state.events[0]).toMatchObject({
notificationStatus: 'failed',
notification: {
status: 'failed',
errorCode: 'SEND_CAPABILITY_UNAVAILABLE'
}
})
expect(mocks.sender.send).not.toHaveBeenCalled()
})
it('does not send twice when the same snapshot diff is checked again', async () => {
installGroupDb()
mocks.capability.getPersonalWechatSendCapability.mockResolvedValue({
ready: true,
capabilities: { text: true, image: false, voice: false },
supported: true,
status: 'ready',
senderStatus: {},
message: 'ready'
})
mocks.sender.send.mockResolvedValue({ success: true })
mocks.chat.getGroupSnapshotAsync
.mockResolvedValueOnce({
roomId: 'room@chatroom',
members: [member, { wxid: 'wxid_other', wechatNickname: '另一个人' }]
})
.mockResolvedValueOnce({ roomId: 'room@chatroom', members: [member] })
.mockResolvedValueOnce({ roomId: 'room@chatroom', members: [member] })
const service = new GroupExitMonitorService()
service.setMonitoredRoomIds(['room@chatroom'], ['room@chatroom'])
await service.start(true)
await service.checkNow()
await service.checkNow()
expect(mocks.sender.send).toHaveBeenCalledOnce()
expect(service.getState().events).toHaveLength(1)
})
})
@@ -645,6 +645,71 @@ describe('scheduled report scheduling', () => {
})
})
it('routes scheduled report sends and retries through the Action Gateway adapter', async () => {
const storageDir = await mkdtemp(join(tmpdir(), 'tracememo-scheduled-report-gateway-'))
const send = vi.fn()
const sendAction = vi
.fn()
.mockResolvedValueOnce({
actionId: 'action-1',
status: 'failed' as const,
decision: 'allow' as const,
errorCode: 'SEND_FAILED' as const,
reason: '连接器暂时不可用',
startedAt: '2026-08-27T01:00:00.000Z',
finishedAt: '2026-08-27T01:00:01.000Z'
})
.mockResolvedValueOnce({
actionId: 'action-2',
status: 'sent' as const,
decision: 'allow' as const,
startedAt: '2026-08-27T01:01:00.000Z',
finishedAt: '2026-08-27T01:01:01.000Z'
})
const service = new ScheduledReportService({
storageDir,
...makeDependencies({ send, sendAction })
})
const created = await service.createTask({
name: 'Gateway 日报',
group: '研发群',
target: '研发群@chatroom',
scheduleTime: '09:00'
})
const first = await runScheduled(service, created.data!.id)
const recovered = await service.retryScheduledReportSend(first.id)
expect(first).toMatchObject({
status: 'partial_success',
pngPath: '/tmp/saved.png',
sendTarget: '研发群@chatroom'
})
expect(recovered.data).toMatchObject({ status: 'success', retryCount: 1 })
expect(send).not.toHaveBeenCalled()
expect(sendAction).toHaveBeenNthCalledWith(
1,
expect.objectContaining({
target: '研发群@chatroom',
filePath: '/tmp/saved.png',
triggerType: 'scheduled',
executionId: first.id,
taskId: created.data!.id
})
)
expect(sendAction).toHaveBeenNthCalledWith(
2,
expect.objectContaining({
target: '研发群@chatroom',
filePath: '/tmp/saved.png',
triggerType: 'scheduled',
executionId: first.id,
retryCount: 1,
taskId: created.data!.id
})
)
})
it('retries an existing PNG without generating the report again', async () => {
const storageDir = await mkdtemp(join(tmpdir(), 'tracememo-scheduled-report-retry-'))
await enableNotifications(storageDir)
+220
View File
@@ -0,0 +1,220 @@
import { mkdtempSync, readJsonSync, rmSync } from 'fs-extra'
import { tmpdir } from 'os'
import { join } from 'path'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
const mocks = vi.hoisted(() => ({
capability: { getPersonalWechatSendCapability: vi.fn() },
sender: { send: vi.fn() }
}))
vi.mock('electron', () => ({ app: { getPath: () => '/tmp/tracememo-gateway-default' } }))
vi.mock('../../src/main/services/personal-wechat-capability-service', () => ({
personalWechatCapabilityService: mocks.capability
}))
vi.mock('../../src/main/services/personal-wechat-send-service', () => ({
personalWechatSendService: mocks.sender
}))
import { WechatActionGateway } from '../../src/main/services/wechat-action-gateway'
import type { PersonalWechatSendCapability } from '../../src/shared/personal-wechat'
import type { WechatActionResult } from '../../src/shared/wechat-action'
const readyCapability: PersonalWechatSendCapability = {
supported: true,
ready: true,
status: 'ready',
capabilities: { text: true, image: true, voice: true },
senderStatus: {} as never,
message: 'ready'
}
describe('WechatActionGateway', () => {
const directories: string[] = []
beforeEach(() => {
mocks.capability.getPersonalWechatSendCapability.mockReset().mockResolvedValue(readyCapability)
mocks.sender.send.mockReset().mockResolvedValue({ success: true, status: {} })
})
afterEach(() => {
while (directories.length) rmSync(directories.pop()!, { recursive: true, force: true })
})
function createGateway(): WechatActionGateway {
const userData = mkdtempSync(join(tmpdir(), 'tracememo-wechat-action-'))
directories.push(userData)
return new WechatActionGateway({ getUserDataPath: () => userData })
}
function memberAction(
gateway: WechatActionGateway,
roomId = 'room@chatroom'
): Promise<WechatActionResult> {
gateway.registerMemberEvent({ id: 'event-1', roomId: 'room@chatroom' })
return gateway.execute({
origin: 'member_monitor',
purpose: 'member_left_notification',
triggerType: 'automation',
sourceId: 'event-1',
recipient: { type: 'group', id: roomId },
content: { type: 'text', text: '张三已退出群聊' }
})
}
it('allows a valid member action, sends once, and writes an audit record', async () => {
const userData = mkdtempSync(join(tmpdir(), 'tracememo-wechat-action-'))
directories.push(userData)
const gateway = new WechatActionGateway({ getUserDataPath: () => userData })
const result = await memberAction(gateway)
expect(result).toMatchObject({ status: 'sent', decision: 'allow' })
expect(mocks.sender.send).toHaveBeenCalledOnce()
expect(mocks.sender.send).toHaveBeenCalledWith({
type: 'text',
to: 'room@chatroom',
isGroup: true,
text: '张三已退出群聊'
})
const audit = readJsonSync(join(userData, 'actions', 'wechat-actions.json'))
expect(audit).toEqual([
expect.objectContaining({
purpose: 'member_left_notification',
recipientId: 'room@chatroom',
sendStatus: 'sent',
decision: 'allow',
contentPreview: '张三已退出群聊'
})
])
})
it('deduplicates the same member event across repeated execution calls', async () => {
const gateway = createGateway()
const first = await memberAction(gateway)
const second = await memberAction(gateway)
expect(second.actionId).toBe(first.actionId)
expect(mocks.sender.send).toHaveBeenCalledOnce()
})
it('blocks a recipient outside the source group', async () => {
const gateway = createGateway()
const result = await memberAction(gateway, 'other@chatroom')
expect(result).toMatchObject({
status: 'blocked',
decision: 'block',
errorCode: 'RECIPIENT_SCOPE_VIOLATION'
})
expect(mocks.sender.send).not.toHaveBeenCalled()
})
it('blocks unknown automation purposes before checking capability', async () => {
const gateway = createGateway()
const result = await gateway.execute({
origin: 'unknown',
purpose: 'arbitrary_message',
triggerType: 'automation',
recipient: { type: 'group', id: 'room@chatroom' },
content: { type: 'text', text: '不应自动发送' }
})
expect(result).toMatchObject({
status: 'blocked',
decision: 'block',
errorCode: 'ACTION_NOT_ALLOWED'
})
expect(mocks.capability.getPersonalWechatSendCapability).not.toHaveBeenCalled()
expect(mocks.sender.send).not.toHaveBeenCalled()
})
it('returns INVALID_REQUEST for malformed input without throwing from audit handling', async () => {
const gateway = createGateway()
const result = await gateway.execute({
origin: 'member_monitor',
purpose: 'member_left_notification',
triggerType: 'automation',
recipient: { type: 'group', id: 'room@chatroom' }
} as never)
expect(result).toMatchObject({
status: 'blocked',
decision: 'block',
errorCode: 'INVALID_REQUEST'
})
})
it('returns INVALID_RECIPIENT when the recipient id is missing', async () => {
const gateway = createGateway()
const result = await gateway.execute({
origin: 'member_monitor',
purpose: 'member_left_notification',
triggerType: 'automation',
sourceId: 'event-1',
recipient: { type: 'group', id: '' },
content: { type: 'text', text: '张三已退出群聊' }
})
expect(result).toMatchObject({
status: 'blocked',
decision: 'block',
errorCode: 'INVALID_RECIPIENT'
})
expect(mocks.sender.send).not.toHaveBeenCalled()
})
it('does not accept an event lookup for a different source id', async () => {
const userData = mkdtempSync(join(tmpdir(), 'tracememo-wechat-action-'))
directories.push(userData)
const gateway = new WechatActionGateway({
getUserDataPath: () => userData,
getMemberEvent: () => ({ id: 'different-event', roomId: 'room@chatroom' })
})
const result = await gateway.execute({
origin: 'member_monitor',
purpose: 'member_left_notification',
triggerType: 'automation',
sourceId: 'event-1',
recipient: { type: 'group', id: 'room@chatroom' },
content: { type: 'text', text: '张三已退出群聊' }
})
expect(result).toMatchObject({
status: 'blocked',
decision: 'block',
errorCode: 'INVALID_REQUEST'
})
expect(mocks.sender.send).not.toHaveBeenCalled()
})
it('returns a structured capability failure without sending', async () => {
const gateway = createGateway()
mocks.capability.getPersonalWechatSendCapability.mockResolvedValueOnce({
...readyCapability,
ready: false,
capabilities: { text: false, image: false, voice: false },
message: '当前微信发送能力不可用'
})
const result = await memberAction(gateway)
expect(result).toMatchObject({
status: 'failed',
decision: 'allow',
errorCode: 'SEND_CAPABILITY_UNAVAILABLE'
})
expect(mocks.sender.send).not.toHaveBeenCalled()
})
it('converts a transport throw into SEND_FAILED', async () => {
const gateway = createGateway()
mocks.sender.send.mockRejectedValueOnce(new Error('connector timeout'))
const result = await memberAction(gateway)
expect(result).toMatchObject({
status: 'failed',
decision: 'allow',
errorCode: 'SEND_FAILED',
reason: 'connector timeout'
})
})
})