mirror of
https://wget.la/https://github.com/Wxw-Gu/WechatExplorer
synced 2026-10-07 06:17:07 +08:00
fix: auto retry scheduled reports after send recovery
This commit is contained in:
@@ -46,6 +46,9 @@ const EXECUTIONS_FILE = 'executions.json'
|
|||||||
const NOTIFICATIONS_FILE = 'notifications.json'
|
const NOTIFICATIONS_FILE = 'notifications.json'
|
||||||
const SETTINGS_FILE = 'settings.json'
|
const SETTINGS_FILE = 'settings.json'
|
||||||
const TICK_MS = 15_000
|
const TICK_MS = 15_000
|
||||||
|
const AUTO_RESEND_MAX_ATTEMPTS = 6
|
||||||
|
const AUTO_RESEND_BASE_BACKOFF_MS = 5 * 60_000
|
||||||
|
const AUTO_RESEND_MAX_BACKOFF_MS = 60 * 60_000
|
||||||
const NOTIFICATION_TEST_MESSAGE = `✅ TraceMemo 定时日报通知已开启
|
const NOTIFICATION_TEST_MESSAGE = `✅ TraceMemo 定时日报通知已开启
|
||||||
|
|
||||||
以后定时日报生成或发送出现异常时,
|
以后定时日报生成或发送出现异常时,
|
||||||
@@ -77,28 +80,6 @@ export interface ScheduledReportSendActionInput {
|
|||||||
taskId?: string
|
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 => ({
|
const defaultDependencies = (): ScheduledReportDependencies => ({
|
||||||
getCapability: () => personalWechatCapabilityService.getPersonalWechatSendCapability(),
|
getCapability: () => personalWechatCapabilityService.getPersonalWechatSendCapability(),
|
||||||
generateReport: generateAgentGroupReport,
|
generateReport: generateAgentGroupReport,
|
||||||
@@ -161,6 +142,29 @@ interface ScheduledReportNotificationPayload {
|
|||||||
suggestedAction?: string
|
suggestedAction?: string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export 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,
|
||||||
|
// filehelper 是文件传输助手:唯一的真实发送验收对象,走 contact 而不是群。
|
||||||
|
recipient: { type: input.target === 'filehelper' ? 'contact' : 'group', id: input.target },
|
||||||
|
content: { type: 'image', path: input.filePath },
|
||||||
|
metadata: {
|
||||||
|
taskId: input.taskId,
|
||||||
|
retryCount: input.retryCount || 0
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
export function validateScheduleTime(value: string): boolean {
|
export function validateScheduleTime(value: string): boolean {
|
||||||
return /^(?:[01]\d|2[0-3]):[0-5]\d$/.test(String(value || '').trim())
|
return /^(?:[01]\d|2[0-3]):[0-5]\d$/.test(String(value || '').trim())
|
||||||
}
|
}
|
||||||
@@ -498,11 +502,7 @@ export class ScheduledReportService {
|
|||||||
await this.load()
|
await this.load()
|
||||||
const execution = this.executions!.find((item) => item.id === executionId)
|
const execution = this.executions!.find((item) => item.id === executionId)
|
||||||
if (!execution) return { success: false, error: '未找到定时日报执行记录' }
|
if (!execution) return { success: false, error: '未找到定时日报执行记录' }
|
||||||
const existing = this.retrying.get(executionId)
|
const result = await this.startRetry(execution)
|
||||||
if (existing) return { success: true, data: await existing }
|
|
||||||
const promise = this.retrySend(execution).finally(() => this.retrying.delete(executionId))
|
|
||||||
this.retrying.set(executionId, promise)
|
|
||||||
const result = await promise
|
|
||||||
return {
|
return {
|
||||||
success: result.status !== 'failed',
|
success: result.status !== 'failed',
|
||||||
data: result,
|
data: result,
|
||||||
@@ -568,6 +568,7 @@ export class ScheduledReportService {
|
|||||||
async tick(at = this.deps.now?.() || new Date()): Promise<void> {
|
async tick(at = this.deps.now?.() || new Date()): Promise<void> {
|
||||||
await this.load()
|
await this.load()
|
||||||
await this.flushNotifications()
|
await this.flushNotifications()
|
||||||
|
await this.flushPendingSends(at)
|
||||||
if (!this.deps.isDatabaseReady()) return
|
if (!this.deps.isDatabaseReady()) return
|
||||||
const nowMs = at.getTime()
|
const nowMs = at.getTime()
|
||||||
for (const task of [...this.tasks!]) {
|
for (const task of [...this.tasks!]) {
|
||||||
@@ -959,6 +960,15 @@ export class ScheduledReportService {
|
|||||||
return completed
|
return completed
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** 为同一执行记录复用唯一的重发 Promise,合并自动重发和手动重试。 */
|
||||||
|
private startRetry(execution: ScheduledReportExecution): Promise<ScheduledReportExecution> {
|
||||||
|
const existing = this.retrying.get(execution.id)
|
||||||
|
if (existing) return existing
|
||||||
|
const promise = this.retrySend(execution).finally(() => this.retrying.delete(execution.id))
|
||||||
|
this.retrying.set(execution.id, promise)
|
||||||
|
return promise
|
||||||
|
}
|
||||||
|
|
||||||
private clearFailureFields(
|
private clearFailureFields(
|
||||||
execution: ScheduledReportExecution,
|
execution: ScheduledReportExecution,
|
||||||
patch: Partial<ScheduledReportExecution>
|
patch: Partial<ScheduledReportExecution>
|
||||||
@@ -1190,6 +1200,40 @@ export class ScheduledReportService {
|
|||||||
await Promise.all([this.saveNotifications(), this.saveExecutions()])
|
await Promise.all([this.saveNotifications(), this.saveExecutions()])
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** waiting_to_send 执行的下一次自动重发时间:自 finishedAt 起按指数退避。 */
|
||||||
|
private autoResendDueAt(execution: ScheduledReportExecution): number {
|
||||||
|
const base = Date.parse(execution.finishedAt || execution.startedAt)
|
||||||
|
const attempt = Math.max(execution.retryCount || 0, 0)
|
||||||
|
return base + Math.min(AUTO_RESEND_BASE_BACKOFF_MS * 2 ** attempt, AUTO_RESEND_MAX_BACKOFF_MS)
|
||||||
|
}
|
||||||
|
|
||||||
|
/** 能力恢复后自动重发挂起的日报;重试有上限,超限后仍可手动重试。 */
|
||||||
|
private async flushPendingSends(at: Date): Promise<void> {
|
||||||
|
await this.load()
|
||||||
|
const nowMs = at.getTime()
|
||||||
|
const candidates = this.executions!.filter(
|
||||||
|
(item) =>
|
||||||
|
item.status === 'waiting_to_send' &&
|
||||||
|
item.retryable !== false &&
|
||||||
|
Boolean(item.pngPath) &&
|
||||||
|
(item.retryCount || 0) < AUTO_RESEND_MAX_ATTEMPTS &&
|
||||||
|
!this.retrying.has(item.id) &&
|
||||||
|
this.autoResendDueAt(item) <= nowMs
|
||||||
|
)
|
||||||
|
const retries: Promise<ScheduledReportExecution>[] = []
|
||||||
|
for (const execution of candidates) {
|
||||||
|
const task = this.tasks!.find((item) => item.id === execution.taskId)
|
||||||
|
if (!task || !task.enabled) continue
|
||||||
|
retries.push(this.startRetry(execution))
|
||||||
|
}
|
||||||
|
const results = await Promise.allSettled(retries)
|
||||||
|
for (const result of results) {
|
||||||
|
if (result.status === 'rejected') {
|
||||||
|
console.warn('[ScheduledReport] auto resend failed:', result.reason)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private notificationCapabilityReasonForSend(
|
private notificationCapabilityReasonForSend(
|
||||||
result: AgentHubNotificationResult
|
result: AgentHubNotificationResult
|
||||||
): ScheduledReportNotificationCapabilityReason {
|
): ScheduledReportNotificationCapabilityReason {
|
||||||
@@ -1219,6 +1263,8 @@ export class ScheduledReportService {
|
|||||||
private resolveTarget(task: ScheduledReportTask): string | undefined {
|
private resolveTarget(task: ScheduledReportTask): string | undefined {
|
||||||
const explicit = String(task.target || '').trim()
|
const explicit = String(task.target || '').trim()
|
||||||
if (explicit.endsWith('@chatroom')) return explicit
|
if (explicit.endsWith('@chatroom')) return explicit
|
||||||
|
// 文件传输助手是唯一允许真实发送验收的个人目标。
|
||||||
|
if (explicit.toLowerCase() === 'filehelper') return 'filehelper'
|
||||||
const contact = resolveMd5(task.group || explicit)
|
const contact = resolveMd5(task.group || explicit)
|
||||||
return contact?.m_nsUsrName?.endsWith('@chatroom') ? contact.m_nsUsrName : undefined
|
return contact?.m_nsUsrName?.endsWith('@chatroom') ? contact.m_nsUsrName : undefined
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import { join } from 'path'
|
|||||||
import { describe, expect, it, vi } from 'vitest'
|
import { describe, expect, it, vi } from 'vitest'
|
||||||
vi.mock('electron', () => ({ app: { getPath: () => '/tmp/tracememo-test-user-data' } }))
|
vi.mock('electron', () => ({ app: { getPath: () => '/tmp/tracememo-test-user-data' } }))
|
||||||
import {
|
import {
|
||||||
|
buildScheduledReportActionRequest,
|
||||||
calculateNextRunAt,
|
calculateNextRunAt,
|
||||||
ScheduledReportService,
|
ScheduledReportService,
|
||||||
validateScheduleTime
|
validateScheduleTime
|
||||||
@@ -973,3 +974,410 @@ describe('scheduled report scheduling', () => {
|
|||||||
).toHaveLength(1)
|
).toHaveLength(1)
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
describe('scheduled report auto resend', () => {
|
||||||
|
const latestExecution = async (
|
||||||
|
service: ScheduledReportService,
|
||||||
|
taskId: string
|
||||||
|
): Promise<ScheduledReportExecution> => (await service.listExecutions(taskId))[0]
|
||||||
|
|
||||||
|
const waitForAutoResend = async (
|
||||||
|
predicate: () => Promise<boolean>,
|
||||||
|
timeoutMs = 3_000
|
||||||
|
): Promise<void> => {
|
||||||
|
const deadline = Date.now() + timeoutMs
|
||||||
|
while (Date.now() < deadline) {
|
||||||
|
if (await predicate()) return
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 10))
|
||||||
|
}
|
||||||
|
throw new Error('auto resend state did not settle in time')
|
||||||
|
}
|
||||||
|
|
||||||
|
const sentAction = {
|
||||||
|
actionId: 'send-success-action',
|
||||||
|
status: 'sent' as const,
|
||||||
|
decision: 'allow' as const,
|
||||||
|
startedAt: '2026-08-27T01:06:00.000Z',
|
||||||
|
finishedAt: '2026-08-27T01:06:01.000Z'
|
||||||
|
}
|
||||||
|
|
||||||
|
const unavailableSendAction = () => async () => ({
|
||||||
|
actionId: 'send-unavailable-action',
|
||||||
|
status: 'failed' as const,
|
||||||
|
decision: 'allow' as const,
|
||||||
|
errorCode: 'SEND_CAPABILITY_UNAVAILABLE' as const,
|
||||||
|
reason: '当前微信发送能力不可用',
|
||||||
|
startedAt: '2026-08-27T01:00:00.000Z',
|
||||||
|
finishedAt: '2026-08-27T01:00:01.000Z'
|
||||||
|
})
|
||||||
|
|
||||||
|
it('auto resends a waiting report after backoff and notifies recovery', async () => {
|
||||||
|
const storageDir = await mkdtemp(join(tmpdir(), 'tracememo-scheduled-report-auto-resend-'))
|
||||||
|
await enableNotifications(storageDir)
|
||||||
|
const clock = { now: Date.parse('2026-08-27T01:00:00.000Z') }
|
||||||
|
const sendNotification = vi.fn(async () => ({
|
||||||
|
success: true as const,
|
||||||
|
status: 'sent' as const
|
||||||
|
}))
|
||||||
|
const sendAction = vi
|
||||||
|
.fn()
|
||||||
|
.mockResolvedValueOnce({
|
||||||
|
actionId: 'send-1',
|
||||||
|
status: 'failed' as const,
|
||||||
|
decision: 'allow' as const,
|
||||||
|
errorCode: 'SEND_CAPABILITY_UNAVAILABLE' as const,
|
||||||
|
reason: '当前微信发送能力不可用',
|
||||||
|
startedAt: '2026-08-27T01:00:00.000Z',
|
||||||
|
finishedAt: '2026-08-27T01:00:01.000Z'
|
||||||
|
})
|
||||||
|
.mockResolvedValueOnce({
|
||||||
|
actionId: 'send-2',
|
||||||
|
status: 'sent' as const,
|
||||||
|
decision: 'allow' as const,
|
||||||
|
startedAt: '2026-08-27T01:06:00.000Z',
|
||||||
|
finishedAt: '2026-08-27T01:06:01.000Z'
|
||||||
|
})
|
||||||
|
const service = new ScheduledReportService({
|
||||||
|
storageDir,
|
||||||
|
now: () => new Date(clock.now),
|
||||||
|
...makeDependencies({ sendAction, sendNotification })
|
||||||
|
})
|
||||||
|
const created = await service.createTask({
|
||||||
|
name: '自动重发日报',
|
||||||
|
group: '研发群',
|
||||||
|
target: '研发群@chatroom',
|
||||||
|
scheduleTime: '09:00'
|
||||||
|
})
|
||||||
|
|
||||||
|
const first = await runScheduled(service, created.data!.id)
|
||||||
|
expect(first).toMatchObject({
|
||||||
|
status: 'waiting_to_send',
|
||||||
|
retryCount: 0,
|
||||||
|
pngPath: '/tmp/saved.png',
|
||||||
|
notificationStatus: 'not_needed'
|
||||||
|
})
|
||||||
|
|
||||||
|
// retryCount 0 的退避是 5 分钟:推进到到期时刻再 tick。
|
||||||
|
clock.now = Date.parse(first.finishedAt || first.startedAt) + 5 * 60_000 + 1_000
|
||||||
|
await service.tick(new Date(clock.now))
|
||||||
|
await waitForAutoResend(async () => sendAction.mock.calls.length >= 2)
|
||||||
|
|
||||||
|
const recovered = await latestExecution(service, created.data!.id)
|
||||||
|
expect(recovered).toMatchObject({
|
||||||
|
status: 'success',
|
||||||
|
retryCount: 1,
|
||||||
|
sendStatus: 'success',
|
||||||
|
notificationStatus: 'sent'
|
||||||
|
})
|
||||||
|
expect(sendAction).toHaveBeenNthCalledWith(
|
||||||
|
2,
|
||||||
|
expect.objectContaining({
|
||||||
|
target: '研发群@chatroom',
|
||||||
|
filePath: '/tmp/saved.png',
|
||||||
|
executionId: first.id,
|
||||||
|
triggerType: 'scheduled',
|
||||||
|
retryCount: 1,
|
||||||
|
taskId: created.data!.id
|
||||||
|
})
|
||||||
|
)
|
||||||
|
expect(sendNotification).toHaveBeenCalledTimes(1)
|
||||||
|
expect(await service.listNotifications()).toEqual([
|
||||||
|
expect.objectContaining({
|
||||||
|
type: 'recovery',
|
||||||
|
status: 'sent',
|
||||||
|
title: expect.stringContaining('已恢复')
|
||||||
|
})
|
||||||
|
])
|
||||||
|
})
|
||||||
|
|
||||||
|
it('does not auto resend before the backoff elapses', async () => {
|
||||||
|
const storageDir = await mkdtemp(join(tmpdir(), 'tracememo-scheduled-report-resend-early-'))
|
||||||
|
await enableNotifications(storageDir)
|
||||||
|
const clock = { now: Date.parse('2026-08-27T01:00:00.000Z') }
|
||||||
|
const sendAction = vi.fn(unavailableSendAction())
|
||||||
|
const sendNotification = vi.fn()
|
||||||
|
const service = new ScheduledReportService({
|
||||||
|
storageDir,
|
||||||
|
now: () => new Date(clock.now),
|
||||||
|
...makeDependencies({ sendAction, sendNotification })
|
||||||
|
})
|
||||||
|
const created = await service.createTask({
|
||||||
|
name: '退避未到期日报',
|
||||||
|
group: '研发群',
|
||||||
|
target: '研发群@chatroom',
|
||||||
|
scheduleTime: '09:00'
|
||||||
|
})
|
||||||
|
const first = await runScheduled(service, created.data!.id)
|
||||||
|
expect(first).toMatchObject({ status: 'waiting_to_send', retryCount: 0 })
|
||||||
|
|
||||||
|
// 仅推进 1 分钟,尚未达到 5 分钟退避。
|
||||||
|
clock.now = Date.parse(first.finishedAt || first.startedAt) + 60_000
|
||||||
|
await service.tick(new Date(clock.now))
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 50))
|
||||||
|
|
||||||
|
expect(sendAction).toHaveBeenCalledTimes(1)
|
||||||
|
expect(await latestExecution(service, created.data!.id)).toMatchObject({
|
||||||
|
status: 'waiting_to_send',
|
||||||
|
retryCount: 0
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
it('stops auto resending after the attempt limit but keeps manual retry', async () => {
|
||||||
|
const storageDir = await mkdtemp(join(tmpdir(), 'tracememo-scheduled-report-resend-limit-'))
|
||||||
|
await enableNotifications(storageDir)
|
||||||
|
const clock = { now: Date.parse('2026-08-27T01:00:00.000Z') }
|
||||||
|
const sendAction = vi.fn(unavailableSendAction())
|
||||||
|
const service = new ScheduledReportService({
|
||||||
|
storageDir,
|
||||||
|
now: () => new Date(clock.now),
|
||||||
|
...makeDependencies({ sendAction })
|
||||||
|
})
|
||||||
|
const created = await service.createTask({
|
||||||
|
name: '重发上限日报',
|
||||||
|
group: '研发群',
|
||||||
|
target: '研发群@chatroom',
|
||||||
|
scheduleTime: '09:00'
|
||||||
|
})
|
||||||
|
let current = await runScheduled(service, created.data!.id)
|
||||||
|
expect(current).toMatchObject({ status: 'waiting_to_send', retryCount: 0 })
|
||||||
|
|
||||||
|
for (let attempt = 1; attempt <= 6; attempt += 1) {
|
||||||
|
const backoff = Math.min(5 * 60_000 * 2 ** (current.retryCount || 0), 60 * 60_000)
|
||||||
|
clock.now = Date.parse(current.finishedAt || current.startedAt) + backoff + 1_000
|
||||||
|
await service.tick(new Date(clock.now))
|
||||||
|
await waitForAutoResend(async () => {
|
||||||
|
const execution = await latestExecution(service, created.data!.id)
|
||||||
|
return sendAction.mock.calls.length === attempt + 1 && execution.retryCount === attempt
|
||||||
|
})
|
||||||
|
// 等待后台 promise 从 retrying map 移除,避免影响下一次 tick 判定。
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 20))
|
||||||
|
current = await latestExecution(service, created.data!.id)
|
||||||
|
}
|
||||||
|
|
||||||
|
// 已达 6 次自动重试上限:继续 tick 不再自动发送。
|
||||||
|
clock.now += 60 * 60_000 + 1_000
|
||||||
|
await service.tick(new Date(clock.now))
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 50))
|
||||||
|
expect(sendAction).toHaveBeenCalledTimes(7)
|
||||||
|
expect(await latestExecution(service, created.data!.id)).toMatchObject({
|
||||||
|
status: 'waiting_to_send',
|
||||||
|
retryCount: 6
|
||||||
|
})
|
||||||
|
|
||||||
|
// 手动重试不受自动上限限制。
|
||||||
|
const manual = await service.retryScheduledReportSend(current.id)
|
||||||
|
expect(manual.success).toBe(true)
|
||||||
|
expect(sendAction).toHaveBeenCalledTimes(8)
|
||||||
|
expect(manual.data).toMatchObject({ status: 'waiting_to_send', retryCount: 7 })
|
||||||
|
})
|
||||||
|
|
||||||
|
it('skips auto resend when the task is disabled', async () => {
|
||||||
|
const storageDir = await mkdtemp(join(tmpdir(), 'tracememo-scheduled-report-resend-disabled-'))
|
||||||
|
await enableNotifications(storageDir)
|
||||||
|
const clock = { now: Date.parse('2026-08-27T01:00:00.000Z') }
|
||||||
|
const sendAction = vi.fn(unavailableSendAction())
|
||||||
|
const service = new ScheduledReportService({
|
||||||
|
storageDir,
|
||||||
|
now: () => new Date(clock.now),
|
||||||
|
...makeDependencies({ sendAction })
|
||||||
|
})
|
||||||
|
const created = await service.createTask({
|
||||||
|
name: '任务禁用日报',
|
||||||
|
group: '研发群',
|
||||||
|
target: '研发群@chatroom',
|
||||||
|
scheduleTime: '09:00'
|
||||||
|
})
|
||||||
|
const first = await runScheduled(service, created.data!.id)
|
||||||
|
expect(first).toMatchObject({ status: 'waiting_to_send', retryCount: 0 })
|
||||||
|
|
||||||
|
await service.setTaskEnabled(created.data!.id, false)
|
||||||
|
clock.now = Date.parse(first.finishedAt || first.startedAt) + 10 * 60_000
|
||||||
|
await service.tick(new Date(clock.now))
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 50))
|
||||||
|
|
||||||
|
expect(sendAction).toHaveBeenCalledTimes(1)
|
||||||
|
expect(await latestExecution(service, created.data!.id)).toMatchObject({
|
||||||
|
status: 'waiting_to_send',
|
||||||
|
retryCount: 0
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
it('sends an expired execution only once when ticks overlap', async () => {
|
||||||
|
const storageDir = await mkdtemp(join(tmpdir(), 'tracememo-scheduled-report-resend-overlap-'))
|
||||||
|
const clock = { now: Date.parse('2026-08-27T01:00:00.000Z') }
|
||||||
|
let releaseRetry: ((value: typeof sentAction) => void) | undefined
|
||||||
|
const retryPending = new Promise<typeof sentAction>((resolve) => {
|
||||||
|
releaseRetry = resolve
|
||||||
|
})
|
||||||
|
const sendAction = vi
|
||||||
|
.fn()
|
||||||
|
.mockResolvedValueOnce(unavailableSendAction()())
|
||||||
|
.mockReturnValueOnce(retryPending)
|
||||||
|
const service = new ScheduledReportService({
|
||||||
|
storageDir,
|
||||||
|
now: () => new Date(clock.now),
|
||||||
|
...makeDependencies({ sendAction })
|
||||||
|
})
|
||||||
|
const created = await service.createTask({
|
||||||
|
name: '并发重发日报',
|
||||||
|
group: '研发群',
|
||||||
|
target: '研发群@chatroom',
|
||||||
|
scheduleTime: '09:00'
|
||||||
|
})
|
||||||
|
const first = await runScheduled(service, created.data!.id)
|
||||||
|
const dueAt = Date.parse(first.finishedAt || first.startedAt) + 5 * 60_000 + 1_000
|
||||||
|
|
||||||
|
const firstTick = service.tick(new Date(dueAt))
|
||||||
|
const secondTick = service.tick(new Date(dueAt))
|
||||||
|
await waitForAutoResend(async () => sendAction.mock.calls.length === 2)
|
||||||
|
expect(sendAction).toHaveBeenCalledTimes(2)
|
||||||
|
|
||||||
|
releaseRetry!(sentAction)
|
||||||
|
await Promise.all([firstTick, secondTick])
|
||||||
|
expect(sendAction).toHaveBeenCalledTimes(2)
|
||||||
|
expect(await latestExecution(service, created.data!.id)).toMatchObject({
|
||||||
|
status: 'success',
|
||||||
|
retryCount: 1
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
it('shares the automatic resend with a concurrent manual retry', async () => {
|
||||||
|
const storageDir = await mkdtemp(
|
||||||
|
join(tmpdir(), 'tracememo-scheduled-report-resend-manual-race-')
|
||||||
|
)
|
||||||
|
const clock = { now: Date.parse('2026-08-27T01:00:00.000Z') }
|
||||||
|
let releaseRetry: ((value: typeof sentAction) => void) | undefined
|
||||||
|
const retryPending = new Promise<typeof sentAction>((resolve) => {
|
||||||
|
releaseRetry = resolve
|
||||||
|
})
|
||||||
|
const sendAction = vi
|
||||||
|
.fn()
|
||||||
|
.mockResolvedValueOnce(unavailableSendAction()())
|
||||||
|
.mockReturnValueOnce(retryPending)
|
||||||
|
const service = new ScheduledReportService({
|
||||||
|
storageDir,
|
||||||
|
now: () => new Date(clock.now),
|
||||||
|
...makeDependencies({ sendAction })
|
||||||
|
})
|
||||||
|
const created = await service.createTask({
|
||||||
|
name: '手动竞争日报',
|
||||||
|
group: '研发群',
|
||||||
|
target: '研发群@chatroom',
|
||||||
|
scheduleTime: '09:00'
|
||||||
|
})
|
||||||
|
const first = await runScheduled(service, created.data!.id)
|
||||||
|
const dueAt = Date.parse(first.finishedAt || first.startedAt) + 5 * 60_000 + 1_000
|
||||||
|
|
||||||
|
const automatic = service.tick(new Date(dueAt))
|
||||||
|
await waitForAutoResend(async () => sendAction.mock.calls.length === 2)
|
||||||
|
const manual = service.retryScheduledReportSend(first.id)
|
||||||
|
expect(sendAction).toHaveBeenCalledTimes(2)
|
||||||
|
|
||||||
|
releaseRetry!(sentAction)
|
||||||
|
await automatic
|
||||||
|
const manualResult = await manual
|
||||||
|
expect(sendAction).toHaveBeenCalledTimes(2)
|
||||||
|
expect(manualResult).toMatchObject({
|
||||||
|
success: true,
|
||||||
|
data: { status: 'success', retryCount: 1 }
|
||||||
|
})
|
||||||
|
expect(await latestExecution(service, created.data!.id)).toMatchObject(manualResult.data)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('resumes an expired waiting execution after service recreation', async () => {
|
||||||
|
const storageDir = await mkdtemp(join(tmpdir(), 'tracememo-scheduled-report-resend-restart-'))
|
||||||
|
const clock = { now: Date.parse('2026-08-27T01:00:00.000Z') }
|
||||||
|
const firstSendAction = vi.fn(unavailableSendAction())
|
||||||
|
const firstService = new ScheduledReportService({
|
||||||
|
storageDir,
|
||||||
|
now: () => new Date(clock.now),
|
||||||
|
...makeDependencies({ sendAction: firstSendAction })
|
||||||
|
})
|
||||||
|
const created = await firstService.createTask({
|
||||||
|
name: '重启恢复日报',
|
||||||
|
group: '研发群',
|
||||||
|
target: '研发群@chatroom',
|
||||||
|
scheduleTime: '09:00'
|
||||||
|
})
|
||||||
|
const first = await runScheduled(firstService, created.data!.id)
|
||||||
|
expect(first.status).toBe('waiting_to_send')
|
||||||
|
|
||||||
|
const resumedSendAction = vi.fn().mockResolvedValue(sentAction)
|
||||||
|
const secondService = new ScheduledReportService({
|
||||||
|
storageDir,
|
||||||
|
now: () => new Date(clock.now),
|
||||||
|
...makeDependencies({ sendAction: resumedSendAction })
|
||||||
|
})
|
||||||
|
const dueAt = Date.parse(first.finishedAt || first.startedAt) + 5 * 60_000 + 1_000
|
||||||
|
await secondService.tick(new Date(dueAt))
|
||||||
|
|
||||||
|
expect(resumedSendAction).toHaveBeenCalledTimes(1)
|
||||||
|
expect(await latestExecution(secondService, created.data!.id)).toMatchObject({
|
||||||
|
status: 'success',
|
||||||
|
retryCount: 1,
|
||||||
|
sendStatus: 'success'
|
||||||
|
})
|
||||||
|
const restored = new ScheduledReportService({ storageDir })
|
||||||
|
expect(await latestExecution(restored, created.data!.id)).toMatchObject({
|
||||||
|
status: 'success',
|
||||||
|
retryCount: 1
|
||||||
|
})
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
||||||
|
describe('scheduled report filehelper target', () => {
|
||||||
|
it('builds contact recipients for filehelper and group recipients for chatrooms', () => {
|
||||||
|
expect(
|
||||||
|
buildScheduledReportActionRequest({
|
||||||
|
target: 'filehelper',
|
||||||
|
filePath: '/tmp/report.png',
|
||||||
|
executionId: 'exec-1',
|
||||||
|
triggerType: 'scheduled',
|
||||||
|
retryCount: 1,
|
||||||
|
taskId: 'task-1'
|
||||||
|
})
|
||||||
|
).toMatchObject({ recipient: { type: 'contact', id: 'filehelper' } })
|
||||||
|
expect(
|
||||||
|
buildScheduledReportActionRequest({
|
||||||
|
target: '12345@chatroom',
|
||||||
|
filePath: '/tmp/report.png',
|
||||||
|
executionId: 'exec-2',
|
||||||
|
triggerType: 'scheduled'
|
||||||
|
})
|
||||||
|
).toMatchObject({ recipient: { type: 'group', id: '12345@chatroom' } })
|
||||||
|
})
|
||||||
|
|
||||||
|
it('sends a waiting report to filehelper when the target is the transfer assistant', async () => {
|
||||||
|
const storageDir = await mkdtemp(join(tmpdir(), 'tracememo-scheduled-report-filehelper-'))
|
||||||
|
const sendAction = vi.fn(async () => ({
|
||||||
|
actionId: 'send-filehelper',
|
||||||
|
status: 'failed' as const,
|
||||||
|
decision: 'allow' as const,
|
||||||
|
errorCode: 'SEND_CAPABILITY_UNAVAILABLE' as const,
|
||||||
|
reason: '当前微信发送能力不可用',
|
||||||
|
startedAt: '2026-08-27T01:00:00.000Z',
|
||||||
|
finishedAt: '2026-08-27T01:00:01.000Z'
|
||||||
|
}))
|
||||||
|
const service = new ScheduledReportService({
|
||||||
|
storageDir,
|
||||||
|
...makeDependencies({ sendAction })
|
||||||
|
})
|
||||||
|
const created = await service.createTask({
|
||||||
|
name: '文件传输助手日报',
|
||||||
|
group: '文件传输助手',
|
||||||
|
target: 'filehelper',
|
||||||
|
scheduleTime: '09:00'
|
||||||
|
})
|
||||||
|
|
||||||
|
const execution = await runScheduled(service, created.data!.id)
|
||||||
|
|
||||||
|
expect(execution).toMatchObject({
|
||||||
|
status: 'waiting_to_send',
|
||||||
|
sendTarget: 'filehelper',
|
||||||
|
pngPath: '/tmp/saved.png'
|
||||||
|
})
|
||||||
|
expect(sendAction).toHaveBeenCalledWith(
|
||||||
|
expect.objectContaining({ target: 'filehelper', filePath: '/tmp/saved.png' })
|
||||||
|
)
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|||||||
Reference in New Issue
Block a user