diff --git a/src/main/services/group-exit-monitor-service.ts b/src/main/services/group-exit-monitor-service.ts index 7d294ef..bfe44b4 100644 --- a/src/main/services/group-exit-monitor-service.ts +++ b/src/main/services/group-exit-monitor-service.ts @@ -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 & + Partial< + Pick + > + class GroupExitMonitorService { + private readonly actionGateway: GroupExitActionGateway private active = false private nativeMonitorActive = false private snapshots = new Map() @@ -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[] = [] 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 { + 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 + 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 diff --git a/src/main/services/scheduled-report-service.ts b/src/main/services/scheduled-report-service.ts index c730778..68282b5 100644 --- a/src/main/services/scheduled-report-service.ts +++ b/src/main/services/scheduled-report-service.ts @@ -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 - send: (request: PersonalWechatSendRequest) => ReturnType + send: (request: PersonalWechatSendRequest) => Promise + sendAction?: (input: ScheduledReportSendActionInput) => Promise sendNotification: (input: { to?: string; text: string }) => Promise 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 { + 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(['today', 'yesterday', '7days', 'recent24h']) const reportRangeLabel = (range: ScheduledReportRange): string => @@ -148,7 +212,21 @@ export class ScheduledReportService { private readonly retrying = new Map>() constructor(deps?: Partial) { - 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 { @@ -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> + 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 { diff --git a/src/main/services/wechat-action-gateway.ts b/src/main/services/wechat-action-gateway.ts new file mode 100644 index 0000000..d0812c8 --- /dev/null +++ b/src/main/services/wechat-action-gateway.ts @@ -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 + send?: (request: PersonalWechatSendRequest) => Promise + getMemberEvent?: ( + sourceId: string + ) => + | WechatActionMemberEventReference + | Promise + | undefined + checkEntitlement?: (request: WechatActionRequest) => boolean | Promise + getUserDataPath?: () => string + now?: () => Date +} + +interface LoadedAuditState { + path: string + records: WechatActionAuditRecord[] +} + +export interface WechatActionPolicyContext { + memberEvent?: WechatActionMemberEventReference +} + +const defaultDependencies = (): Required< + Pick +> => ({ + getCapability: () => personalWechatCapabilityService.getPersonalWechatSendCapability(), + send: (request) => personalWechatSendService.send(request), + getUserDataPath: () => app.getPath('userData'), + now: () => new Date() +}) + +/** + * 统一管理微信发送操作,在真正发送前完成必要检查; + * 具体发送由底层服务处理,使用者只需提供发送内容和对象。 + */ +export class WechatActionGateway { + private readonly deps: Required< + Pick + > & + Omit + private readonly memberEvents = new Map() + private readonly inFlight = new Map>() + 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 { + 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 { + 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 { + 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 { + 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 + 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 { + 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 + 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 diff --git a/src/renderer/src/features/group-exit-monitor/GroupExitMonitorWorkspace.tsx b/src/renderer/src/features/group-exit-monitor/GroupExitMonitorWorkspace.tsx index a61bb50..16e6741 100644 --- a/src/renderer/src/features/group-exit-monitor/GroupExitMonitorWorkspace.tsx +++ b/src/renderer/src/features/group-exit-monitor/GroupExitMonitorWorkspace.tsx @@ -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 ( +
+ {details.label} + {details.description ? {details.description} : null} +
+ ) +} + function SendCapabilityStatus({ capability, className = '' @@ -783,6 +815,7 @@ export function GroupExitMonitorWorkspace({ 群人数 {event.previousCount} 人 → {event.currentCount} 人

) : null} +
微信名
diff --git a/src/renderer/src/styles/group-exit-monitor.scss b/src/renderer/src/styles/group-exit-monitor.scss index da7a142..c97acfc 100644 --- a/src/renderer/src/styles/group-exit-monitor.scss +++ b/src/renderer/src/styles/group-exit-monitor.scss @@ -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)); diff --git a/src/shared/group-exit-monitor.ts b/src/shared/group-exit-monitor.ts index 4e3706a..24c24fd 100644 --- a/src/shared/group-exit-monitor.ts +++ b/src/shared/group-exit-monitor.ts @@ -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 diff --git a/src/shared/wechat-action.ts b/src/shared/wechat-action.ts new file mode 100644 index 0000000..b0d19cb --- /dev/null +++ b/src/shared/wechat-action.ts @@ -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 +} + +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 +} diff --git a/tests/unit/group-exit-monitor-service.test.ts b/tests/unit/group-exit-monitor-service.test.ts index 0f6add4..3ee1b76 100644 --- a/tests/unit/group-exit-monitor-service.test.ts +++ b/tests/unit/group-exit-monitor-service.test.ts @@ -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) + }) }) diff --git a/tests/unit/scheduled-report-service.test.ts b/tests/unit/scheduled-report-service.test.ts index c23de37..1a10af8 100644 --- a/tests/unit/scheduled-report-service.test.ts +++ b/tests/unit/scheduled-report-service.test.ts @@ -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) diff --git a/tests/unit/wechat-action-gateway.test.ts b/tests/unit/wechat-action-gateway.test.ts new file mode 100644 index 0000000..91d44a9 --- /dev/null +++ b/tests/unit/wechat-action-gateway.test.ts @@ -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 { + 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' + }) + }) +})