'use strict'; const assert = require('assert/strict'); const fs = require('fs'); const os = require('os'); const path = require('path'); const { GroupAgentService } = require('../mcp/src/dashboard/group-agent-service'); function output() { return { content: '已收到,我先为您整理后续安排。', confidence: 0.95, intent: '需求确认', reason: '群消息恢复测试固定输出。', requiresHuman: true, citations: [], toolTrace: [], }; } async function main() { const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'qiwei-group-recovery-')); const roomId = 'group-recovery-room'; const sourceMessageId = 'group-recovery-inbound-1'; const messages = [{ msgUniqueIdentifier: sourceMessageId, msgServerId: sourceMessageId, fromRoomId: roomId, msgType: 1, senderId: 'customer-1', senderName: '测试客户', content: '请发一下服务介绍', timestamp: Date.now(), }]; const groups = { [roomId]: { roomName: '恢复测试群' } }; const statePath = path.join(dir, 'group-agent-replies.json'); const runtime = { db: { globalState: () => ({ paused: false }) }, config: { qiwei: { groupGenerationRecoveryRetryCooldownMs: 1000 } }, agent: { async run() { return output(); } }, qiwei: { async sendText() { return { isSendSuccess: true }; } }, }; const options = { projectRoot: dir, statePath, getRuntime: () => runtime, getAccount: () => ({ uid: 'fixture', userId: 'self' }), loadGroups: () => groups, loadMessages: () => messages, appendMessage: (_roomId, message) => messages.push(message), }; try { // Model a runtime exit after the inbound message is archived and the relay // has ACKed it, but before the in-memory generation promise is started. const beforeRestart = new GroupAgentService(options); const job = beforeRestart.queueGenerationJob(roomId, sourceMessageId, 'relay_callback'); assert.equal(job.status, 'queued'); const afterRestart = new GroupAgentService(options); const recovered = await afterRestart.recoverPendingGenerations('smoke:restart'); assert.equal(recovered.recovered, 1); const state = afterRestart.loadState(); const stored = afterRestart.roomState(state, roomId); const persisted = afterRestart.findGenerationJob(stored, sourceMessageId); assert.equal(persisted.status, 'draft_ready'); assert.equal(afterRestart.publicState(roomId).pendingReply.sourceMessageId, sourceMessageId); assert.equal(stored.audit.some(item => item.action === 'group_generation_started'), true); // A second scan must not regenerate or create a duplicate draft. const second = await afterRestart.recoverPendingGenerations('smoke:dedupe'); assert.equal(second.recovered, 0); assert.equal(stored.drafts.length, 1); process.stdout.write(JSON.stringify({ status: 'passed', checks: [ '已归档群消息拥有持久化 generation job', '运行时重启后仅恢复一次并生成待审核草稿', '终态任务不会被重复回放', ] }, null, 2) + '\n'); } finally { fs.rmSync(dir, { recursive: true, force: true }); } } main().catch(error => { process.stderr.write(JSON.stringify({ status: 'failed', message: error.message, stack: error.stack }, null, 2) + '\n'); process.exitCode = 1; });