| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972 |
- const fs = require('fs');
- const path = require('path');
- const crypto = require('crypto');
- const { PACKAGE_ROOT, WORKSPACE_ROOT, latestPath, categoryDir, outputsRoot } = require('../core/output-paths');
- const { AgentWorkbenchDb } = require('../core/agent-workbench-db');
- const { AgentKnowledgeStore } = require('../core/agent-knowledge');
- const { QiweiAgentRuntime, extractExplicitCustomerIntelligence } = require('../core/agent-runtime');
- const { getCustomerSessionGuide } = require('../core/agent-session-guide');
- const { AgentWorkbenchService } = require('../core/agent-workbench-service');
- const { friendlyAgentError } = require('../core/agent-error-message');
- const { VoiceCloneService } = require('../core/voice-clone-service');
- const {
- searchTodoUsers,
- createTodoKnowledge,
- completeTodoKnowledge,
- } = require('./official-office-knowledge-service');
- const { createCustomerTaskOfficialSync } = require('../core/customer-task-official-sync');
- const { messageTimestamp, roomIdOf, isGroupMessage, messageContent, evaluatePolledMessage, normalizePersonalIntakeMode } = require('../core/agent-poller-policy');
- const { sortConversationsByRecency, maxConversationTimestamp } = require('../core/conversation-order');
- const { saveQiweiClientConfig, setActiveQiweiContext, readFmodeVoiceToken } = require('../core/credentials');
- const { FmodeQiweiClient } = require('../providers/fmode-agent-transport');
- const { responseMonitor } = require('./response-monitor-service');
- const { normalizeAllowlistIds, normalizeAllowlistContact } = require('../core/allowlist-config');
- const { getProductMode } = require('../core/product-mode');
- const { getRelayDaemonStatus } = require('../core/relay-daemon');
- const PROJECT_ROOT = WORKSPACE_ROOT;
- const ENV_FILE = path.join(PROJECT_ROOT, '.env.local');
- const INTAKE_AUTOPILOT_CONFIRMATION = 'ENABLE_AUTO_ENROLL_AUTOPILOT';
- const INTAKE_AUTOPILOT_CONFIRMATION_VERSION = 'v1';
- const AUTOPILOT_CONFIRMATION = 'ENABLE_AUTOPILOT';
- function requireAutopilotConfirmation(mode, confirmation, scope = 'global') {
- if (mode !== 'autopilot' || confirmation === AUTOPILOT_CONFIRMATION) return;
- throw new Error(scope === 'conversation' ? '开启会话全自动接管需要二次确认' : '开启全自动接管需要二次确认');
- }
- function readEnvFile(filePath) {
- try {
- const env = {};
- for (const rawLine of fs.readFileSync(filePath, 'utf8').replace(/^\uFEFF/, '').split(/\r?\n/)) {
- const line = rawLine.trim();
- if (!line || line.startsWith('#') || !line.includes('=')) continue;
- const index = line.indexOf('=');
- const key = line.slice(0, index).trim();
- let value = line.slice(index + 1).trim();
- if ((value.startsWith('"') && value.endsWith('"')) || (value.startsWith("'") && value.endsWith("'"))) value = value.slice(1, -1);
- env[key] = value;
- }
- return env;
- } catch {
- return {};
- }
- }
- function refreshAllowedSendersFromEnv(config = {}, envFile = ENV_FILE) {
- if (config.usePersistedAllowlist === true) {
- const combined = normalizeAllowlistIds([...(config.manualAllowedSenders || []), ...(config.autoEnrolledSenders || [])]);
- const current = Array.isArray(config.allowedSenders) ? config.allowedSenders.map(String) : [];
- const changed = combined.length !== current.length || combined.some((id, index) => id !== current[index]);
- if (changed) config.allowedSenders = combined;
- return { changed, count: combined.length };
- }
- const env = readEnvFile(envFile);
- if (!Object.prototype.hasOwnProperty.call(env, 'QIWEI_AUTO_REPLY_ALLOWED_SENDERS')) {
- return { changed: false, count: Array.isArray(config.allowedSenders) ? config.allowedSenders.length : 0 };
- }
- const next = normalizeAllowlistIds(env.QIWEI_AUTO_REPLY_ALLOWED_SENDERS);
- const combined = normalizeAllowlistIds([...next, ...(config.autoEnrolledSenders || [])]);
- const current = Array.isArray(config.allowedSenders) ? config.allowedSenders.map(String) : [];
- const changed = combined.length !== current.length || combined.some((id, index) => id !== current[index]);
- config.manualAllowedSenders = next;
- if (changed) config.allowedSenders = combined;
- return { changed, count: combined.length };
- }
- function storedAllowlistIds(db, key, fallback = []) {
- const raw = db.getSetting(key, '');
- if (!raw) {
- const ids = normalizeAllowlistIds(fallback);
- db.setSetting(key, JSON.stringify(ids));
- return ids;
- }
- try { return normalizeAllowlistIds(JSON.parse(raw)); } catch { return []; }
- }
- function hydrateAccountAllowlist(config, db) {
- const manual = storedAllowlistIds(db, 'manual_allowed_sender_ids', config.qiwei.manualAllowedSenders || config.qiwei.allowedSenders || []);
- const automatic = db.getAutoEnrolledContactIds();
- config.qiwei.manualAllowedSenders = manual;
- config.qiwei.autoEnrolledSenders = automatic;
- config.qiwei.allowedSenders = normalizeAllowlistIds([...manual, ...automatic]);
- config.qiwei.usePersistedAllowlist = true;
- return config.qiwei.allowedSenders;
- }
- function readRuntimeState() {
- try {
- const filePath = path.join(outputsRoot(), 'runtime', 'qiwei-runtime.json');
- return JSON.parse(fs.readFileSync(filePath, 'utf8').replace(/^\uFEFF/, ''));
- } catch {
- return {};
- }
- }
- function runtimeProcessAlive(pid) {
- const numericPid = Number(pid);
- if (!Number.isInteger(numericPid) || numericPid <= 0) return false;
- try {
- process.kill(numericPid, 0);
- return true;
- } catch {
- return false;
- }
- }
- function listenerStateWithRuntime(product, localState) {
- const runtime = readRuntimeState();
- if (!runtimeProcessAlive(runtime.pid) || !['starting', 'running'].includes(runtime.status)) return localState;
- const component = product.mode === 'enterprise'
- ? runtime.components?.enterpriseRelay
- : runtime.components?.personalPolling;
- if (!component) return localState;
- return {
- ...localState,
- running: Boolean(localState.running || ['starting', 'running'].includes(component.status)),
- lastError: component.lastError || localState.lastError || '',
- transport: runtime.transport,
- runtimePid: runtime.pid,
- runtimeStatus: component.status,
- };
- }
- function readClaudeSettingsEnv() {
- const result = {};
- const home = process.env.USERPROFILE || process.env.HOME || '';
- for (const filePath of [path.join(home, '.claude', 'settings.json'), path.join(home, '.claude', 'settings.local.json')]) {
- try {
- const parsed = JSON.parse(fs.readFileSync(filePath, 'utf8').replace(/^\uFEFF/, ''));
- for (const [key, item] of Object.entries(parsed.env || {})) {
- if (!result[key] && typeof item === 'string' && item.trim()) result[key] = item.trim();
- }
- } catch {}
- }
- return result;
- }
- const fileEnv = readEnvFile(ENV_FILE);
- const claudeEnv = readClaudeSettingsEnv();
- function value(name, fallback = '') {
- const candidates = [process.env[name], fileEnv[name], claudeEnv[name], fallback];
- return String(candidates.find(item => typeof item === 'string' && item.trim()) ?? '').trim();
- }
- function bool(name, fallback = false) {
- return /^(1|true|yes|on)$/i.test(value(name, fallback ? 'true' : 'false'));
- }
- function number(name, fallback, min = -Infinity, max = Infinity) {
- const parsed = Number(value(name, String(fallback)));
- return Math.min(max, Math.max(min, Number.isFinite(parsed) ? parsed : fallback));
- }
- function resolvePath(input, fallback) {
- const selected = input || fallback;
- return selected ? path.resolve(PROJECT_ROOT, selected) : '';
- }
- function loadAgentConfig(overrides = {}) {
- const provider = value('AGENT_PROVIDER', 'claude-code');
- const anthropic = provider === 'anthropic';
- const claudeCode = provider === 'claude-code';
- const manualAllowedSenders = normalizeAllowlistIds(value('QIWEI_AUTO_REPLY_ALLOWED_SENDERS'));
- const baseConfig = {
- accountKey: accountRuntimeKey({ uid: value('QIWEI_UID') || value('QIWE_UID'), guid: value('QIWEI_GUID') || value('QIWE_GUID') }),
- dbPath: resolvePath(value('QIWEI_AGENT_DB_PATH'), latestPath('messages', 'agent-workbench.db')),
- legacyDbPath: path.resolve(PACKAGE_ROOT, '..', '..', 'qiwei-agent-workbench', 'data', 'workbench.db'),
- globalDefaultPaused: bool('QIWEI_AGENT_GLOBAL_DEFAULT_PAUSED', false),
- conversationDefaultMode: value('QIWEI_AGENT_DEFAULT_MODE', 'review'),
- autoSendConfidence: number('QIWEI_AGENT_AUTO_SEND_CONFIDENCE', 0.88, 0, 1),
- intake: {
- mode: normalizePersonalIntakeMode(value('QIWEI_PERSONAL_INTAKE_MODE', 'allowlist_only')),
- welcomeEnabled: bool('QIWEI_WELCOME_ENABLED', false),
- welcomeText: value('QIWEI_WELCOME_TEXT', '您好,已经收到您的消息,我会尽快为您处理。'),
- welcomeSendMode: value('QIWEI_WELCOME_SEND_MODE', 'draft') === 'send' ? 'send' : 'draft',
- effectiveAutopilotMode: 'autopilot',
- },
- knowledgeDir: resolvePath(
- value('QIWEI_AGENT_KNOWLEDGE_DIR'),
- fs.existsSync(path.join(PROJECT_ROOT, 'knowledge')) ? path.join(PROJECT_ROOT, 'knowledge') : path.join(PACKAGE_ROOT, 'knowledge')
- ),
- contextFiles: value('QIWEI_AGENT_CONTEXT_FILES', 'personality.md,context.md,rules.md').split(',').map(item => item.trim()).filter(Boolean),
- contextCharLimit: number('QIWEI_AGENT_CONTEXT_CHAR_LIMIT', 6000, 1000, 20000),
- memory: {
- enabled: bool('QIWEI_AGENT_MEMORY_ENABLED', true),
- coreCharLimit: number('QIWEI_AGENT_MEMORY_CORE_CHAR_LIMIT', 4000, 500, 8000),
- recallLimit: number('QIWEI_AGENT_MEMORY_RECALL_LIMIT', 6, 0, 20),
- recentMessageLimit: number('QIWEI_AGENT_MEMORY_RECENT_MESSAGES', 12, 4, 40),
- historyScanLimit: number('QIWEI_AGENT_MEMORY_HISTORY_SCAN_LIMIT', 240, 20, 1000),
- extractionIntervalMs: number('QIWEI_AGENT_MEMORY_EXTRACTION_INTERVAL_MS', 1000, 250, 60000),
- extractionRetryBaseMs: number('QIWEI_AGENT_MEMORY_EXTRACTION_RETRY_BASE_MS', 1000, 250, 60000),
- extractionMaxAttempts: number('QIWEI_AGENT_MEMORY_EXTRACTION_MAX_ATTEMPTS', 3, 1, 10),
- },
- agent: {
- provider,
- apiKey: value('AGENT_API_KEY') || value(anthropic ? 'ANTHROPIC_AUTH_TOKEN' : 'OPENAI_API_KEY') || value('QIWEI_AUTO_REPLY_AI_KEY'),
- baseUrl: (value('AGENT_BASE_URL') || value(anthropic ? 'ANTHROPIC_BASE_URL' : 'OPENAI_BASE_URL') || value('QIWEI_AUTO_REPLY_AI_BASE_URL') || (anthropic ? 'https://api.anthropic.com' : 'https://api.openai.com/v1')).replace(/\/$/, ''),
- model: value('AGENT_MODEL') || value(anthropic || claudeCode ? 'ANTHROPIC_MODEL' : 'OPENAI_MODEL') || value('QIWEI_AUTO_REPLY_AI_MODEL') || (anthropic || claudeCode ? 'sonnet' : 'gpt-4.1-mini'),
- maxToolRounds: number('AGENT_MAX_TOOL_ROUNDS', 4, 1, 8),
- promptCharLimit: number('QIWEI_AGENT_PROMPT_CHAR_LIMIT', 16000, 2000, 50000),
- systemPromptCharLimit: number('QIWEI_AGENT_SYSTEM_PROMPT_CHAR_LIMIT', 12000, 2000, 50000),
- qualityPassScore: number('QIWEI_AGENT_QUALITY_PASS_SCORE', 82, 60, 100),
- businessGoal: value('QIWEI_AGENT_BUSINESS_GOAL'),
- stagePlaybook: value('QIWEI_AGENT_STAGE_PLAYBOOK'),
- claudeExecutable: value('CLAUDE_CODE_EXECUTABLE') || path.join(path.dirname(process.execPath), 'node_modules', '@anthropic-ai', 'claude-code', 'bin', 'claude.exe'),
- claudeWorkdir: resolvePath(value('CLAUDE_CODE_WORKDIR'), PROJECT_ROOT),
- claudeSessionFile: resolvePath(value('CLAUDE_CODE_SESSION_FILE'), latestPath('messages', 'claude-code-sessions.json')),
- claudeProjectId: value('QIWEI_AGENT_PROJECT_ID') || crypto.createHash('sha256').update(PROJECT_ROOT).digest('hex').slice(0, 16),
- claudeMainSessionId: value('QIWEI_AGENT_MAIN_SESSION_ID'),
- claudeTimeoutMs: number('CLAUDE_CODE_TIMEOUT_MS', 120000, 15000, 300000),
- claudeMaxBudgetUsd: number('CLAUDE_CODE_MAX_BUDGET_USD', 1, 0.05, 5),
- claudeRetryMaxBudgetUsd: number('CLAUDE_CODE_RETRY_MAX_BUDGET_USD', 3, 0.1, 5),
- claudeBare: bool('CLAUDE_CODE_BARE', true),
- claudeEffort: value('CLAUDE_CODE_EFFORT', 'low'),
- claudeTools: value('CLAUDE_CODE_ALLOWED_TOOLS', 'Read,Glob,Grep'),
- claudeSessionMaxTurns: number('CLAUDE_CODE_SESSION_MAX_TURNS', 12, 1, 100),
- claudeSessionMaxAgeMs: number('CLAUDE_CODE_SESSION_MAX_AGE_MS', 86400000, 60000, 604800000),
- },
- qiwei: {
- transport: 'fmode-gateway',
- authToken: value('QIWEI_AUTH_TOKEN'),
- uid: value('QIWEI_UID') || value('QIWE_UID'),
- guid: value('QIWEI_GUID') || value('QIWE_GUID'),
- apiBase: value('QIWEI_API_BASE') || value('QIWE_API_BASE'),
- userId: '',
- nickname: '',
- corpName: '',
- allowedSenders: [...manualAllowedSenders],
- manualAllowedSenders: [...manualAllowedSenders],
- autoEnrolledSenders: [],
- accountKey: '',
- intakeAutopilotConversationMode: 'autopilot',
- selfUserId: value('QIWEI_AUTO_REPLY_SELF_USER_ID'),
- intervalMs: number('QIWEI_AUTO_REPLY_INTERVAL_MS', 10000, 3000, 60000),
- initialSyncLimit: number('QIWEI_AGENT_INITIAL_SYNC_LIMIT', 5000, 100, 5000),
- initialSyncMaxPages: number('QIWEI_AGENT_INITIAL_SYNC_MAX_PAGES', 200, 10, 500),
- startupGraceSeconds: number('QIWEI_AGENT_STARTUP_GRACE_SECONDS', 10, 0, 60),
- trustConfiguredOnStatusError: bool('QIWEI_TRUST_CONFIGURED_ON_STATUS_ERROR', false),
- responseReminderMinutes: number('QIWEI_RESPONSE_REMINDER_MINUTES', 15, 1, 1440),
- responseUrgentMinutes: number('QIWEI_RESPONSE_URGENT_MINUTES', 60, 1, 10080),
- },
- voice: {
- endpoint: value('QIWEI_VOICE_ENDPOINT', 'https://server.fmode.cn/api/voice/indextts2'),
- authToken: readFmodeVoiceToken({
- voiceAuthToken: value('QIWEI_VOICE_AUTH_TOKEN'),
- endpoint: value('QIWEI_VOICE_ENDPOINT', 'https://server.fmode.cn/api/voice/indextts2'),
- }),
- model: 'fmode-voice',
- requestTimeoutMs: number('QIWEI_TTS_TIMEOUT_MS', 180000, 15000, 300000),
- },
- };
- const config = {
- ...baseConfig,
- ...overrides,
- memory: { ...baseConfig.memory, ...(overrides.memory || {}) },
- intake: { ...baseConfig.intake, ...(overrides.intake || {}) },
- agent: { ...baseConfig.agent, ...(overrides.agent || {}) },
- qiwei: { ...baseConfig.qiwei, ...(overrides.qiwei || {}) },
- voice: { ...baseConfig.voice, ...(overrides.voice || {}) },
- };
- config.qiwei.accountKey = config.accountKey;
- config.qiwei.intakeAutopilotConversationMode = config.intake.effectiveAutopilotMode;
- config.agent.claudeAddDirs = [config.knowledgeDir].filter(Boolean);
- return config;
- }
- function accountRuntimeKey(input = {}) {
- const source = String(input.uid || input.guid || input.userId || 'default').trim();
- return crypto.createHash('sha256').update(source || 'default').digest('hex').slice(0, 16);
- }
- function accountWorkbenchOverrides(input = {}) {
- const storageKey = accountRuntimeKey(input);
- return {
- accountKey: storageKey,
- dbPath: latestPath('messages', `agent-workbench-${storageKey}.db`),
- agent: {
- claudeSessionFile: latestPath('messages', `claude-code-sessions-${storageKey}.json`),
- claudeProjectId: `${crypto.createHash('sha256').update(PROJECT_ROOT).digest('hex').slice(0, 12)}-${storageKey.slice(0, 8)}`,
- },
- qiwei: {
- uid: String(input.uid || '').trim(),
- guid: String(input.guid || '').trim(),
- userId: String(input.userId || '').trim(),
- nickname: String(input.nickname || '').trim(),
- corpName: String(input.corpName || '').trim(),
- apiBase: String(input.apiBase || '').trim(),
- ...(input.useConfiguredAllowlist === true ? {} : { allowedSenders: [], manualAllowedSenders: [] }),
- },
- };
- }
- function backfillCustomerIntelligence(db) {
- const version = '4';
- if (db.getSetting('customer_intelligence_backfill_version', '') === version) return { skipped: true, profileFields: 0, taskCount: 0, alertCount: 0 };
- const removed = db.lastIntelligenceMigration || { tasks: { removed: 0 }, alerts: { removed: 0 } };
- let profileFields = 0;
- let taskCount = 0;
- let alertCount = 0;
- for (const conversation of db.listConversations()) {
- const existing = db.getProfile(conversation.id);
- let current = { profile: { ...existing.profile }, tags: existing.tags || [] };
- for (const message of db.listMessages(conversation.id, 200).filter(item => item.direction === 'inbound')) {
- const intelligence = extractExplicitCustomerIntelligence(message.content, current.profile || {}, {});
- if (Object.keys(intelligence.profileUpdates).length) {
- current = { profile: { ...current.profile, ...intelligence.profileUpdates }, tags: current.tags };
- profileFields += Object.keys(intelligence.profileUpdates).length;
- }
- const changed = Object.keys(intelligence.profileUpdates).length > 0;
- const tasks = intelligence.tasks.map(item => ({ ...item, sourceMessageId: changed ? message.id : null }));
- const alerts = intelligence.alerts.map(item => ({ ...item, sourceMessageId: changed || item.managedBy === 'event' ? message.id : null }));
- taskCount += db.reconcileCustomerTasks(conversation.id, tasks, message.id).tasks.length;
- alertCount += db.reconcileCustomerAlerts(conversation.id, alerts, message.id).alerts.length;
- }
- if (Object.keys(current.profile).length) db.updateProfile(conversation.id, { ...existing.profile, ...current.profile }, existing.tags);
- }
- db.setSetting('customer_intelligence_backfill_version', version);
- if (profileFields || taskCount || alertCount) {
- db.audit({ actor: 'migration', action: 'customer_intelligence_backfilled', detail: { version, removed, profileFields, taskCount, alertCount } });
- }
- return { version, removed, profileFields, taskCount, alertCount };
- }
- function backfillCustomerMemory(service) {
- const version = '2';
- const db = service?.db;
- if (!db || !service.memory || db.getSetting('customer_memory_backfill_version', '') === version) {
- return { skipped: true, conversations: 0, captured: 0 };
- }
- let conversations = 0;
- let captured = 0;
- for (const conversation of db.listConversations()) {
- const result = service.memory.backfillConversation(conversation);
- conversations += 1;
- captured += Number(result.captured || 0);
- }
- db.setSetting('customer_memory_backfill_version', version);
- db.audit({ actor: 'migration', action: 'customer_memory_backfilled', detail: { version, conversations, captured } });
- return { version, conversations, captured };
- }
- const delay = ms => new Promise(resolve => setTimeout(resolve, ms));
- class QiweiAgentPoller {
- constructor({ config, db, qiwei, service }) {
- this.config = config;
- this.db = db;
- this.qiwei = qiwei;
- this.service = service;
- this.running = false;
- this.loopPromise = null;
- this.startedAt = 0;
- this.lastError = '';
- this.wakeLoop = null;
- }
- status() {
- return {
- running: this.running,
- syncKey: Number(this.db.getPollState('sync_key', '0')),
- lastError: this.lastError,
- startedAt: this.startedAt || null,
- };
- }
- async start() {
- if (this.running) return this.status();
- if (!this.qiwei.isConfigured()) throw new Error('Fmode 企微网关尚未配置,请先完成 Fmode 鉴权和企微扫码登录');
- if (this.config.allowedSenders.length === 0 && this.db.intakePolicy().mode === 'allowlist_only') throw new Error('企微联系人白名单为空,拒绝启动');
- await this.qiwei.checkLogin();
- this.running = true;
- this.startedAt = Math.floor(Date.now() / 1000);
- this.lastError = '';
- let syncKey = Number(this.db.getPollState('sync_key', '0')) || 0;
- if (syncKey === 0) syncKey = await this.establishBaseline();
- this.db.audit({ actor: 'runtime', action: 'poller_started', detail: { syncKey, allowlistCount: this.config.allowedSenders.length } });
- this.loopPromise = this.loop(syncKey);
- return this.status();
- }
- stop() {
- if (!this.running) return this.status();
- this.running = false;
- if (this.wakeLoop) this.wakeLoop();
- this.db.audit({ actor: 'runtime', action: 'poller_stopped', detail: this.status() });
- return this.status();
- }
- waitInterval() {
- return new Promise(resolve => {
- const timer = setTimeout(() => {
- this.wakeLoop = null;
- resolve();
- }, this.config.intervalMs);
- this.wakeLoop = () => {
- clearTimeout(timer);
- this.wakeLoop = null;
- resolve();
- };
- });
- }
- async establishBaseline() {
- let cursor = 0;
- let reachedEnd = false;
- let pages = 0;
- let count = 0;
- while (pages < this.config.initialSyncMaxPages) {
- const result = await this.qiwei.syncMessages(cursor, this.config.initialSyncLimit);
- const list = result.syncMsgList || [];
- pages += 1;
- count += list.length;
- const seqs = list.map(item => Number(item.seq)).filter(Number.isFinite);
- const next = Math.max(cursor, Number(result.travelSyncKey) || 0, seqs.length ? Math.max(...seqs) : 0);
- if (list.length === 0) { reachedEnd = true; break; }
- if (next <= cursor) throw new Error(`历史消息游标没有前进(seq=${cursor})`);
- cursor = next;
- }
- if (!reachedEnd) throw new Error(`历史消息超过 ${this.config.initialSyncMaxPages} 页,拒绝启动自动处理`);
- this.db.setPollState('sync_key', cursor);
- this.db.audit({ actor: 'runtime', action: 'poller_baseline_established', detail: { cursor, pages, skippedHistory: count } });
- return cursor;
- }
- async loop(initialSyncKey) {
- let syncKey = initialSyncKey;
- while (this.running) {
- try {
- const result = await this.qiwei.syncMessages(syncKey, 50);
- const list = result.syncMsgList || [];
- const seqs = list.map(item => Number(item.seq)).filter(Number.isFinite);
- const next = Math.max(syncKey, Number(result.travelSyncKey) || 0, seqs.length ? Math.max(...seqs) : 0);
- for (const message of list) await this.process(message);
- syncKey = next;
- this.db.setPollState('sync_key', syncKey);
- this.lastError = '';
- } catch (error) {
- this.lastError = error.message;
- this.db.audit({ actor: 'runtime', action: 'poller_error', detail: { message: error.message } });
- }
- if (this.running) await this.waitInterval();
- }
- }
- async process(message) {
- const allowlist = refreshAllowedSendersFromEnv(this.config);
- if (allowlist.changed) {
- this.db.audit({ actor: 'runtime', action: 'poller_allowlist_reloaded', detail: { count: allowlist.count } });
- }
- return ingestMessageForWorkbench({ config: this.config, db: this.db, service: this.service, qiwei: this.qiwei }, message, 'polling');
- }
- }
- const OUTBOUND_NOTIFICATION_MSG_TYPE = 2001;
- const OUTBOUND_CONTENT_MSG_TYPES = new Set([0, 1, 2]);
- function outboundContactId(message = {}, config = {}) {
- const senderId = String(message.senderId || '');
- const receiverId = String(message.receiverId || '');
- const allowedSenders = new Set((config.allowedSenders || []).map(String));
- const selfUserIds = new Set([
- String(config.selfUserId || ''),
- String(config.userId || ''),
- String(message.userId || ''),
- ].filter(Boolean));
- return selfUserIds.has(senderId) && allowedSenders.has(receiverId) ? receiverId : '';
- }
- function recordOutboundMessage(target, message, source, contactId) {
- const content = messageContent(message);
- if (!content) throw new Error('Outbound message content is empty');
- const timestamp = messageTimestamp(message.timestamp);
- if (!timestamp) throw new Error('Outbound message timestamp is invalid');
- const existing = target.db.getConversationByContactId(contactId);
- const conversation = target.db.ensureConversation(contactId, existing?.contact_name || message.receiverName || '白名单企微联系人');
- const inserted = target.db.insertMessage({
- conversationId: conversation.id,
- externalId: String(message.msgServerId || message.msgUniqueIdentifier || `${message.senderId}:${message.seq}`),
- direction: 'outbound',
- senderType: 'human',
- content,
- status: 'sent',
- createdAt: new Date(timestamp * 1000).toISOString(),
- raw: { ...message, source },
- });
- if (inserted.created) {
- target.db.audit({
- actor: 'runtime',
- action: 'outbound_message_captured',
- conversationId: conversation.id,
- entityId: inserted.message.id,
- detail: { source, seq: Number(message.seq) || 0 },
- });
- target.service?.emit?.('change', { type: 'message', conversationId: conversation.id });
- }
- return { status: inserted.created ? 'outbound_recorded' : 'outbound_duplicate', message: inserted.message };
- }
- async function syncOutboundNotification(target, notification, source, contactId) {
- if (!target.qiwei?.syncMessages) throw new Error('Outbound message sync is unavailable');
- const notificationSeq = Number(notification.seq) || 0;
- const senderId = String(notification.senderId || '');
- const delays = [0, 250, 750, 1500];
- for (const delayMs of delays) {
- if (delayMs) await new Promise(resolve => setTimeout(resolve, delayMs));
- const result = await target.qiwei.syncMessages(Math.max(0, notificationSeq - 1), 10);
- const message = (result.syncMsgList || []).find(item => (
- Number(item.seq) >= notificationSeq
- && String(item.senderId || '') === senderId
- && String(item.receiverId || '') === contactId
- && OUTBOUND_CONTENT_MSG_TYPES.has(Number(item.msgType))
- && Boolean(messageContent(item))
- ));
- if (message) return recordOutboundMessage(target, message, `${source}:outbound-sync`, contactId);
- }
- throw new Error(`Outbound message content is not available after notification seq=${notificationSeq}`);
- }
- async function ingestMessageForWorkbench(target, message = {}, source = 'callback') {
- if (isGroupMessage(message)) {
- return { status: 'ignored_group_reply_disabled' };
- }
- const contactId = outboundContactId(message, target.config);
- if (contactId && Number(message.msgType) === OUTBOUND_NOTIFICATION_MSG_TYPE) {
- return syncOutboundNotification(target, message, source, contactId);
- }
- if (contactId && OUTBOUND_CONTENT_MSG_TYPES.has(Number(message.msgType)) && messageContent(message)) {
- return recordOutboundMessage(target, message, source, contactId);
- }
- const intakePolicy = target.db.intakePolicy();
- const candidate = evaluatePolledMessage(message, { ...target.config, intakeMode: intakePolicy.mode });
- if (!candidate.eligible) return { status: `ignored_${candidate.reason || 'ineligible'}` };
- const { content, senderId, timestamp } = candidate;
- let onboarding = false;
- let allowAutoSend;
- if (candidate.unknownContact) {
- if (intakePolicy.mode === 'allowlist_only') return { status: 'ignored_not_allowlisted' };
- if (intakePolicy.mode === 'auto_enroll_autopilot' && intakePolicy.autopilotConfirmation !== INTAKE_AUTOPILOT_CONFIRMATION_VERSION) {
- target.db.audit({ actor: 'policy', action: 'auto_enroll_blocked_unconfirmed', detail: { contactId: senderId, source } });
- return { status: 'ignored_auto_enroll_unconfirmed' };
- }
- const autoIds = target.db.addAutoEnrolledContactId(senderId);
- target.config.autoEnrolledSenders = autoIds;
- target.config.allowedSenders = normalizeAllowlistIds([...(target.config.manualAllowedSenders || []), ...autoIds]);
- const conversation = target.db.ensureConversation(senderId, message.senderName || '企微客户');
- const effectiveMode = intakePolicy.mode === 'auto_enroll_review' ? 'review' : 'autopilot';
- target.service.setConversationMode(conversation.id, effectiveMode, 'policy:auto-enroll');
- const welcomeState = !intakePolicy.welcomeEnabled
- ? 'skipped'
- : (intakePolicy.mode === 'auto_enroll_autopilot' && intakePolicy.welcomeSendMode === 'send' ? 'send_pending' : 'draft_pending');
- const accountKey = String(target.config.accountKey || 'default');
- const existing = target.db.getOnboarding(accountKey, senderId);
- target.db.upsertOnboarding({
- accountKey,
- contactId: senderId,
- conversationId: conversation.id,
- contactName: message.senderName || '企微客户',
- source,
- policy: intakePolicy.mode,
- welcomeState,
- welcomeText: intakePolicy.welcomeText,
- });
- if (!existing) {
- target.db.audit({ actor: 'policy', action: 'contact_auto_enrolled', conversationId: conversation.id, detail: { contactId: senderId, source, policy: intakePolicy.mode, effectiveConversationMode: effectiveMode } });
- }
- onboarding = true;
- allowAutoSend = intakePolicy.mode !== 'auto_enroll_review';
- }
- return target.service.ingestInbound({
- externalId: String(message.msgServerId || message.msgUniqueIdentifier || `${senderId}:${message.seq}`),
- contactId: senderId,
- contactName: message.senderName || '企微客户',
- content,
- timestamp: new Date(timestamp * 1000).toISOString(),
- raw: {
- ...message,
- source,
- fromRoomId: message.fromRoomId || null,
- },
- }, { onboarding, allowAutoSend });
- }
- function createWorkbench(overrides = {}) {
- const config = loadAgentConfig(overrides.config || {});
- const db = overrides.db || new AgentWorkbenchDb(config.dbPath, {
- globalPaused: config.globalDefaultPaused,
- defaultMode: config.conversationDefaultMode,
- autoSendConfidence: config.autoSendConfidence,
- personalIntakeMode: config.intake.mode,
- welcomeEnabled: config.intake.welcomeEnabled,
- welcomeText: config.intake.welcomeText,
- welcomeSendMode: config.intake.welcomeSendMode,
- });
- hydrateAccountAllowlist(config, db);
- const staleWelcomeCount = db.markStaleWelcomeSending(
- config.accountKey,
- new Date(Date.now() - 5 * 60 * 1000).toISOString(),
- );
- if (staleWelcomeCount) {
- db.audit({ actor: 'runtime', action: 'onboarding_welcome_delivery_unknown', detail: { count: staleWelcomeCount } });
- }
- if (!overrides.db && !overrides.skipLegacyImport) {
- try { db.importCompatibleDatabase(config.legacyDbPath); } catch (error) {
- db.audit({ actor: 'migration', action: 'legacy_workbench_import_failed', detail: { message: error.message } });
- }
- const removed = db.cleanupInboundContentDuplicates(60);
- if (removed) db.audit({ actor: 'migration', action: 'duplicate_messages_cleaned', detail: { removed } });
- backfillCustomerIntelligence(db);
- }
- const knowledge = overrides.knowledge || new AgentKnowledgeStore({
- knowledgeDir: config.knowledgeDir,
- contextFiles: config.contextFiles,
- contextCharLimit: config.contextCharLimit,
- });
- backfillSentVoiceAudioPaths(db);
- const qiwei = overrides.qiwei || new FmodeQiweiClient(config.qiwei);
- const voice = overrides.voice || new VoiceCloneService({ config: config.voice, qiwei });
- const agent = overrides.agent || new QiweiAgentRuntime({ config: config.agent, knowledge });
- const service = overrides.service || new AgentWorkbenchService({ db, agent, qiwei, config });
- backfillCustomerMemory(service);
- const poller = overrides.poller || new QiweiAgentPoller({ config: config.qiwei, db, qiwei, service });
- return { config, db, knowledge, qiwei, voice, agent, service, poller };
- }
- function configuredStartupAccount() {
- return {
- uid: value('QIWEI_UID') || value('QIWE_UID'),
- guid: value('QIWEI_GUID') || value('QIWE_GUID'),
- apiBase: value('QIWEI_API_BASE') || value('QIWE_API_BASE'),
- userId: '',
- nickname: '',
- corpName: '',
- useConfiguredAllowlist: true,
- };
- }
- const startupAccount = configuredStartupAccount();
- let workbench = createWorkbench(startupAccount.uid ? {
- config: accountWorkbenchOverrides(startupAccount),
- skipLegacyImport: true,
- } : {});
- function migrateStartupWorkbench(target) {
- if (!startupAccount.uid || target.db.listConversations().length) return { imported: false, reason: 'target_not_empty' };
- const candidates = [latestPath('messages', 'agent-workbench.db'), target.config.legacyDbPath];
- for (const sourcePath of candidates) {
- try {
- const result = target.db.importCompatibleDatabase(sourcePath);
- if (!result.imported) continue;
- const removed = target.db.cleanupInboundContentDuplicates(60);
- backfillCustomerIntelligence(target.db);
- target.db.audit({
- actor: 'migration',
- action: 'account_workbench_migration_completed',
- detail: { source: path.basename(sourcePath), removedDuplicates: removed },
- });
- return result;
- } catch (error) {
- target.db.audit({
- actor: 'migration',
- action: 'account_workbench_migration_failed',
- detail: { source: path.basename(sourcePath), message: error.message },
- });
- }
- }
- return { imported: false, reason: 'source_missing_or_incompatible' };
- }
- migrateStartupWorkbench(workbench);
- let accountStatusCache = { checkedAt: 0, value: null };
- let accountStatusRefresh = null;
- let accountStatusRefreshKey = '';
- const accountLastOnlineAt = new Map();
- const accountOfflineChecks = new Map();
- const ONLINE_STATUS_GRACE_MS = 60000;
- const workbenches = new Map();
- function activeAccountMetadata() {
- const context = workbench.qiwei.context();
- return {
- uid: String(context.uid || workbench.config.qiwei.uid || '').trim(),
- guid: String(context.guid || workbench.config.qiwei.guid || '').trim(),
- apiBase: String(context.apiBase || workbench.config.qiwei.apiBase || '').trim(),
- userId: String(workbench.config.qiwei.userId || '').trim(),
- nickname: String(workbench.config.qiwei.nickname || '').trim(),
- corpName: String(workbench.config.qiwei.corpName || '').trim(),
- };
- }
- function applyActiveAccountContext(account) {
- setActiveQiweiContext(account);
- Object.assign(workbench.config.qiwei, account);
- if (workbench.qiwei && workbench.qiwei.config) Object.assign(workbench.qiwei.config, account);
- }
- function provisionalAccountStatus(selected = activeAccountMetadata(), statusText = '正在检测账号状态') {
- return {
- uid: selected.uid,
- guid: selected.guid,
- userId: selected.userId,
- configured: workbench.qiwei.isConfigured(),
- online: false,
- nickname: selected.nickname || selected.userId || '当前企微账号',
- corpName: selected.corpName || '',
- statusCode: null,
- statusText,
- };
- }
- const initialAccount = activeAccountMetadata();
- workbenches.set(accountRuntimeKey(initialAccount), workbench);
- applyActiveAccountContext(initialAccount);
- async function switchActiveAccount(input = {}) {
- const requestedGuid = String(input.guid || '').trim();
- const account = {
- uid: String(input.uid || '').trim(),
- guid: requestedGuid === 'server-managed' ? '' : requestedGuid,
- apiBase: String(input.apiBase || '').trim(),
- userId: String(input.userId || '').trim(),
- nickname: String(input.nickname || input.userId || '').trim(),
- corpName: String(input.corpName || '').trim(),
- };
- if (!account.uid) throw new Error('该账号缺少 Fmode 设备 uid,请重新扫码绑定后再切换');
- const current = activeAccountMetadata();
- const currentKey = accountRuntimeKey(current);
- const nextKey = accountRuntimeKey(account);
- const accountChanged = currentKey !== nextKey;
- if (accountChanged && workbench.poller.status().running) workbench.poller.stop();
- let nextWorkbench = workbenches.get(nextKey);
- if (!nextWorkbench) {
- nextWorkbench = createWorkbench({
- config: accountWorkbenchOverrides(account),
- skipLegacyImport: true,
- });
- workbenches.set(nextKey, nextWorkbench);
- }
- workbench = nextWorkbench;
- const currentAccount = activeAccountMetadata();
- const accountPatch = {
- ...account,
- apiBase: account.apiBase || currentAccount.apiBase,
- guid: account.guid || (accountChanged ? '' : currentAccount.guid),
- };
- applyActiveAccountContext({ ...currentAccount, ...accountPatch });
- saveQiweiClientConfig({
- uid: accountPatch.uid,
- guid: accountPatch.guid,
- apiBase: accountPatch.apiBase,
- envRoot: PROJECT_ROOT,
- });
- const status = provisionalAccountStatus();
- accountStatusCache = { checkedAt: Date.now(), value: status };
- void refreshAccountStatus();
- return {
- status: 'ok',
- assistantMessage: `已切换到账号:${status.nickname || account.nickname || account.userId || account.uid}`,
- summary: {
- switched: accountChanged,
- storageKey: nextKey,
- online: status.online,
- listenerStopped: accountChanged,
- },
- data: { account: status },
- };
- }
- async function refreshAccountStatus() {
- const selected = activeAccountMetadata();
- const targetWorkbench = workbench;
- const refreshKey = accountRuntimeKey(selected);
- if (accountStatusRefresh && accountStatusRefreshKey === refreshKey) return accountStatusRefresh;
- accountStatusRefreshKey = refreshKey;
- accountStatusRefresh = (async () => {
- let next;
- try {
- const data = await targetWorkbench.qiwei.checkLogin();
- const online = Number(data.userOnlineStatus) === 2 && Number(data.errorCode || 0) === 0;
- next = {
- uid: selected.uid,
- guid: selected.guid,
- userId: data.userId || selected.userId,
- configured: data.configured !== false,
- online,
- nickname: data.nickname || selected.nickname || selected.userId || '当前企微账号',
- corpName: data.corpName || selected.corpName || '',
- statusCode: data.userOnlineStatus ?? null,
- statusText: online ? (data.statusFallback ? '账号在线(状态接口待复核)' : '账号在线') : '账号离线',
- };
- if (online) {
- accountLastOnlineAt.set(refreshKey, Date.now());
- accountOfflineChecks.set(refreshKey, 0);
- } else {
- const offlineChecks = Number(accountOfflineChecks.get(refreshKey) || 0) + 1;
- accountOfflineChecks.set(refreshKey, offlineChecks);
- const lastOnlineAt = Number(accountLastOnlineAt.get(refreshKey) || 0);
- if (offlineChecks < 2 && lastOnlineAt && Date.now() - lastOnlineAt < ONLINE_STATUS_GRACE_MS) {
- next.online = true;
- next.statusCode = 2;
- next.statusText = '账号在线(正在复核)';
- }
- }
- } catch {
- next = {
- ...provisionalAccountStatus(selected, '状态检测失败'),
- configured: targetWorkbench.qiwei.isConfigured(),
- };
- const lastOnlineAt = Number(accountLastOnlineAt.get(refreshKey) || 0);
- if (lastOnlineAt && Date.now() - lastOnlineAt < ONLINE_STATUS_GRACE_MS) {
- next.online = true;
- next.statusCode = 2;
- next.statusText = '账号在线(状态刷新中)';
- }
- }
- if (accountRuntimeKey(activeAccountMetadata()) === refreshKey) {
- accountStatusCache = { checkedAt: Date.now(), value: next };
- }
- if (next.userId) {
- targetWorkbench.config.qiwei.userId = String(next.userId);
- targetWorkbench.config.qiwei.selfUserId = String(next.userId);
- }
- return next;
- })().finally(() => {
- if (accountStatusRefreshKey === refreshKey) {
- accountStatusRefresh = null;
- accountStatusRefreshKey = '';
- }
- });
- return accountStatusRefresh;
- }
- async function detectAccountStatus(force = false) {
- const selected = activeAccountMetadata();
- const selectedKey = accountRuntimeKey(selected);
- const cacheMatches = accountStatusCache.value && accountRuntimeKey(accountStatusCache.value) === selectedKey;
- if (force) return refreshAccountStatus();
- if (cacheMatches) {
- if (Date.now() - accountStatusCache.checkedAt >= 8000) void refreshAccountStatus();
- return accountStatusCache.value;
- }
- const provisional = provisionalAccountStatus(selected);
- accountStatusCache = { checkedAt: Date.now(), value: provisional };
- void refreshAccountStatus();
- return provisional;
- }
- function maskedId(value) {
- const text = String(value || '');
- if (text.length <= 4) return '测试联系人';
- return `${text.slice(0, 2)}***${text.slice(-2)}`;
- }
- function completeness(profile = {}) {
- const values = [
- profile.preferredRegion || profile.region || profile.district || profile.intent_area || profile.districts,
- profile.budgetWan || profile.budget || profile.budgetMax || profile.budget_max,
- profile.layout || profile.rooms || profile.house_type,
- profile.area || profile.areaMin || profile.area_min,
- profile.decoration,
- profile.timeline || profile.urgency,
- ];
- return Math.round(values.filter(value => Array.isArray(value) ? value.length : Boolean(value)).length / values.length * 100);
- }
- function parseJson(value, fallback) {
- try { return value ? JSON.parse(value) : fallback; }
- catch { return fallback; }
- }
- function conversationChannelInfo(row, messages = []) {
- const contactId = String(row?.contact_id || '').trim();
- let roomId = '';
- for (const message of messages) {
- const raw = parseJson(message.raw_json, {});
- if (!isGroupMessage(raw)) continue;
- roomId = roomIdOf(raw) || contactId;
- break;
- }
- if (!roomId && /(?:@chatroom$|^(?:room|group|chatroom|r[-_:]))/i.test(contactId)) roomId = contactId;
- const channelType = roomId ? 'group' : 'private';
- return {
- channelType,
- channelLabel: channelType === 'group' ? '客户群聊' : '客户私聊',
- channelId: roomId || contactId,
- };
- }
- function publicConversation(row) {
- const detail = workbench.service.conversationDetail(row.id);
- const claudeSession = getCustomerSessionGuide(row, {
- sessionFile: workbench.config.agent.claudeSessionFile,
- });
- const drafts = detail.drafts || [];
- const latestInbound = [...(detail.messages || [])].reverse().find(message => message.direction === 'inbound') || null;
- const currentDrafts = latestInbound
- ? drafts.filter(item => item.inbound_message_id === latestInbound.id)
- : [];
- const pending = currentDrafts.find(item => item.status === 'pending') || null;
- const latestDraft = pending || currentDrafts.find(item => ['sent', 'approved'].includes(item.status)) || null;
- const currentAgentOutcome = latestInbound && detail.agentOutcome?.entityId === latestInbound.id ? detail.agentOutcome : null;
- const agentError = ['agent_failed', 'agent_not_configured', 'autopilot_send_failed'].includes(currentAgentOutcome?.action)
- ? { ...currentAgentOutcome, ...friendlyAgentError(currentAgentOutcome.message) }
- : null;
- const agentNotice = currentAgentOutcome?.action === 'agent_no_reply_needed' ? currentAgentOutcome : null;
- const rawProfile = detail.profile?.profile || {};
- const { __evidence: profileEvidence = {}, ...profile } = rawProfile;
- const customerTasks = detail.tasks || [];
- const customerAlerts = detail.alerts || [];
- const citations = latestDraft?.citations || [];
- const cutoverAt = Date.parse(workbench.db.getSetting('agent_cutover_at', '')) || Date.now();
- const displayMessages = [...new Map((detail.messages || []).map(message => [message.id, message])).values()]
- .sort((a, b) => Date.parse(a.created_at) - Date.parse(b.created_at))
- .slice(-60);
- const channel = conversationChannelInfo(row, displayMessages);
- const visibleEntityIds = new Set(displayMessages.map(message => message.id));
- const visibleAudit = (detail.audit || []).filter(item =>
- Date.parse(item.created_at) >= cutoverAt || visibleEntityIds.has(item.entity_id)
- ).slice(0, 40).map(item => {
- const publicDetail = { ...(item.detail || {}) };
- delete publicDetail.rawMessage;
- return { ...item, detail_json: undefined, detail: publicDetail };
- });
- return {
- id: row.id,
- displayName: row.contact_name || '白名单测试联系人',
- maskedId: maskedId(row.contact_id),
- ...channel,
- mode: row.mode,
- source: 'live',
- claudeSession,
- messages: displayMessages.map(message => {
- const raw = parseJson(message.raw_json, {});
- if (Number(raw.msgType) === 16 && raw.audioPath) {
- raw.audioAvailable = true;
- delete raw.audioPath;
- }
- return {
- id: message.id,
- role: message.direction === 'inbound' ? 'customer' : message.sender_type,
- content: message.content,
- timestamp: message.created_at,
- status: message.status,
- source: message.direction === 'inbound' ? 'live' : message.sender_type,
- raw,
- };
- }),
- analysis: {
- intent: latestDraft?.intent || '',
- intentLabel: latestDraft?.intent || (agentError ? 'Agent 上游不可用' : agentNotice ? '无需回复' : '待 Agent 处理'),
- demand: profile,
- completenessScore: completeness(profile),
- knowledgeSources: citations.length
- ? citations.map(item => `${item.heading || item.source}${item.source ? ` · ${item.source}` : ''}`)
- : ['真实企微消息', '客户画像', '企业规则库与知识库'],
- reasoning: latestDraft?.reason || agentError?.message || agentNotice?.message || '消息已进入真实企微链路,等待 Agent 生成可审核草稿。',
- },
- customerIntelligence: {
- profile,
- profileEvidence,
- profileUpdatedAt: detail.profile?.updatedAt || null,
- tags: detail.profile?.tags || [],
- tasks: customerTasks.map(item => ({
- id: item.id,
- businessKey: item.business_key,
- type: item.type,
- title: item.title,
- owner: item.owner,
- dueAt: item.due_at,
- priority: item.priority,
- status: item.status,
- reason: item.reason,
- evidence: item.evidence,
- evidenceItems: parseJson(item.evidence_json, []),
- resolutionReason: item.resolution_reason,
- officialTodoId: item.official_todo_id,
- officialSyncStatus: item.official_sync_status,
- officialSyncedAt: item.official_synced_at,
- updatedAt: item.updated_at,
- })),
- alerts: customerAlerts.map(item => ({
- id: item.id,
- businessKey: item.business_key,
- type: item.type,
- severity: item.severity,
- title: item.title,
- detail: item.detail,
- evidence: item.evidence,
- evidenceItems: parseJson(item.evidence_json, []),
- resolutionReason: item.resolution_reason,
- recommendedAction: item.recommended_action,
- status: item.status,
- updatedAt: item.updated_at,
- })),
- memories: (detail.memories || []).map(item => ({
- id: item.id,
- type: item.type,
- content: item.content,
- confidence: item.confidence,
- importance: item.importance,
- status: item.status,
- sourceMessageIds: item.source_message_ids || [],
- direction: item.direction,
- createdBy: item.created_by,
- expiresAt: item.expires_at,
- updatedAt: item.updated_at,
- })),
- memorySnapshot: detail.memorySnapshot ? {
- version: detail.memorySnapshot.version,
- coreChars: String(detail.memorySnapshot.compact_text || '').length,
- createdAt: detail.memorySnapshot.created_at,
- } : null,
- summary: {
- openTasks: customerTasks.filter(item => ['open', 'in_progress'].includes(item.status)).length,
- openAlerts: customerAlerts.filter(item => item.status === 'open').length,
- highAlerts: customerAlerts.filter(item => item.status === 'open' && ['high', 'critical'].includes(item.severity)).length,
- activeMemories: (detail.memories || []).filter(item => item.status === 'active').length,
- },
- },
- pendingReply: pending ? {
- id: pending.id,
- content: pending.content,
- status: pending.status,
- confidence: pending.confidence,
- reason: pending.reason,
- requiresHuman: pending.requires_human,
- citations: pending.citations,
- toolTrace: pending.tool_trace,
- createdAt: pending.created_at,
- } : null,
- drafts,
- audit: visibleAudit,
- agentError,
- agentNotice,
- onboarding: detail.onboarding ? {
- state: detail.onboarding.welcome_state,
- attemptCount: detail.onboarding.attempt_count,
- lastError: detail.onboarding.last_error || '',
- lastAttemptAt: detail.onboarding.last_attempt_at || null,
- sentAt: detail.onboarding.sent_at || null,
- } : null,
- lastMessageAt: maxConversationTimestamp(displayMessages.at(-1)?.created_at, row.last_message_at),
- updatedAt: row.updated_at,
- };
- }
- async function getAgentStatus() {
- const account = await detectAccountStatus();
- const state = workbench.service.state(workbench.poller.status());
- const product = getProductMode();
- const listener = listenerStateWithRuntime(product, state.poller);
- const daemonCallback = getRelayDaemonStatus();
- const runtimeRelayRunning = product.mode === 'enterprise'
- && listener.transport === 'server_relay'
- && listener.running;
- const callback = runtimeRelayRunning
- ? {
- ...daemonCallback,
- running: true,
- processAlive: true,
- pid: listener.runtimePid || daemonCallback.pid,
- heartbeatAt: readRuntimeState().updatedAt || daemonCallback.heartbeatAt,
- }
- : daemonCallback;
- const ingress = product.mode === 'enterprise'
- ? { mode: 'server_callback', running: callback.running, durableQueue: true, ...callback }
- : { mode: 'local_polling', running: state.poller.running, durableQueue: false };
- return {
- status: 'ok',
- data: {
- globalMode: state.global.paused ? 'paused' : state.global.defaultMode,
- global: state.global,
- listener,
- ingress,
- callback,
- account,
- product,
- agent: state.agent,
- knowledge: workbench.knowledge.stats(),
- voice: workbench.voice.status(),
- config: {
- allowedSenderCount: state.qiwei.allowlistCount,
- testMode: false,
- demoMode: false,
- transport: state.qiwei.transport,
- pollIntervalMs: workbench.config.qiwei.intervalMs,
- personalIngress: intakePolicyPayload(workbench),
- },
- safety: {
- whitelistEnabled: state.qiwei.allowlistCount > 0,
- defaultReviewMode: true,
- globalPauseSupported: true,
- messageSendRequiresWhitelist: true,
- demoReplyDisabled: true,
- },
- },
- };
- }
- async function ingestWebhookMessage(message = {}) {
- return ingestMessageForWorkbench({
- config: workbench.config.qiwei,
- db: workbench.db,
- service: workbench.service,
- qiwei: workbench.qiwei,
- }, message, 'relay_callback');
- }
- async function getAllowlistCandidates() {
- const selectedIds = [...workbench.config.qiwei.allowedSenders];
- const candidates = new Map();
- let remoteCount = 0;
- let warning = '';
- try {
- const account = await detectAccountStatus(true);
- if (account.online) {
- const result = await workbench.qiwei.listExternalContacts({ limit: 500, maxPages: 10 });
- remoteCount = Number(result.contactCount) || result.contacts.length;
- for (const item of result.contacts) {
- const contact = normalizeAllowlistContact(item);
- if (contact) candidates.set(contact.id, contact);
- }
- } else {
- warning = '当前企微账号离线,只显示已保存或已有会话中的联系人';
- }
- } catch {
- warning = '联系人列表暂时读取失败,仍可管理已保存联系人或手工添加联系人 ID';
- }
- for (const row of workbench.db.listConversations()) {
- const contact = normalizeAllowlistContact({ contactId: row.contact_id }, row.contact_name);
- if (contact && !candidates.has(contact.id)) candidates.set(contact.id, contact);
- }
- for (const id of selectedIds) {
- if (!candidates.has(id)) candidates.set(id, normalizeAllowlistContact({ id }, '已保存联系人'));
- }
- const contacts = [...candidates.values()]
- .map(contact => ({ ...contact, selected: selectedIds.includes(contact.id) }))
- .sort((a, b) => Number(b.selected) - Number(a.selected) || a.displayName.localeCompare(b.displayName, 'zh-CN'));
- return {
- status: 'ok',
- assistantMessage: warning || `已读取 ${contacts.length} 位可选联系人`,
- summary: { selectedCount: selectedIds.length, candidateCount: contacts.length, remoteCount },
- data: { selectedIds, contacts },
- warnings: warning ? [warning] : [],
- errors: [],
- };
- }
- function updateAllowlistForWorkbench(target, input = {}) {
- refreshAllowedSendersFromEnv(target.config.qiwei);
- const previousIds = new Set(target.config.qiwei.allowedSenders.map(String));
- const ids = normalizeAllowlistIds(input.contactIds || input.allowedSenders || []);
- const retainedAutomatic = target.db.getAutoEnrolledContactIds().filter(id => ids.includes(id));
- target.db.setSetting('manual_allowed_sender_ids', JSON.stringify(ids));
- target.db.setSetting('auto_enrolled_contact_ids', JSON.stringify(retainedAutomatic));
- target.config.qiwei.manualAllowedSenders = [...ids];
- target.config.qiwei.autoEnrolledSenders = retainedAutomatic;
- target.config.qiwei.allowedSenders = normalizeAllowlistIds([...ids, ...retainedAutomatic]);
- let listenerStopped = false;
- const contactNames = input.contactNames || {};
- const addedIds = ids.filter(id => !previousIds.has(id));
- for (const contactId of addedIds) {
- target.db.ensureConversation(contactId, String(contactNames[contactId] || ''));
- }
- if (ids.length && input.autoStart !== false) {
- target.db.setSetting('listener_enabled', 'true');
- target.service.setGlobal({ paused: false }, input.actor || 'allowlist');
- } else if (!ids.length) {
- target.db.setSetting('listener_enabled', 'false');
- }
- if (!ids.length && target.poller.status().running) {
- target.poller.stop();
- listenerStopped = true;
- }
- target.db.audit({
- actor: input.actor || 'human',
- action: 'allowlist_updated',
- detail: { count: ids.length, addedCount: addedIds.length, listenerStopped, autoStart: ids.length > 0 && input.autoStart !== false },
- });
- return {
- status: 'ok',
- assistantMessage: ids.length
- ? `白名单已保存,共 ${ids.length} 位联系人;AI 监听将自动保持开启并按当前审核策略处理`
- : '白名单已清空,AI 监听已停止',
- summary: { selectedCount: ids.length, addedCount: addedIds.length, listenerStopped, autoStart: ids.length > 0 && input.autoStart !== false },
- data: { selectedIds: ids, selectedCount: ids.length, addedIds },
- warnings: [],
- errors: [],
- };
- }
- function addAllowlistContacts(input = {}) {
- refreshAllowedSendersFromEnv(workbench.config.qiwei);
- const contactIds = normalizeAllowlistIds(input.contactIds || []);
- return updateAllowlist({
- ...input,
- contactIds: [...workbench.config.qiwei.allowedSenders, ...contactIds],
- autoStart: input.autoStart !== false,
- actor: input.actor || 'agent:test-contact',
- });
- }
- function intakePolicyPayload(target = workbench) {
- const policy = target.db.intakePolicy();
- const onboardings = target.db.listOnboardings(target.config.accountKey, { limit: 500 });
- return {
- mode: policy.mode,
- welcomeEnabled: policy.welcomeEnabled,
- welcomeText: policy.welcomeText,
- welcomeSendMode: policy.welcomeSendMode,
- effectiveConversationMode: policy.mode === 'auto_enroll_autopilot' ? 'autopilot' : 'review',
- autopilotConfirmed: policy.autopilotConfirmation === INTAKE_AUTOPILOT_CONFIRMATION_VERSION,
- onboarding: {
- total: onboardings.length,
- failed: onboardings.filter(item => item.welcome_state === 'failed').length,
- deliveryUnknown: onboardings.filter(item => item.welcome_state === 'delivery_unknown').length,
- },
- };
- }
- function updateAllowlist(input = {}) {
- return updateAllowlistForWorkbench(workbench, input);
- }
- function getIntakePolicy() {
- return { status: 'ok', data: intakePolicyPayload(workbench), warnings: [], errors: [] };
- }
- function updateIntakePolicyForWorkbench(target, input = {}, actor = 'human') {
- const current = target.db.intakePolicy();
- const mode = input.mode === undefined ? current.mode : String(input.mode);
- const welcomeSendMode = input.welcomeSendMode === undefined ? current.welcomeSendMode : String(input.welcomeSendMode);
- const welcomeEnabled = input.welcomeEnabled === undefined ? current.welcomeEnabled : input.welcomeEnabled === true;
- const welcomeText = input.welcomeText === undefined ? current.welcomeText : String(input.welcomeText || '').trim();
- if (!['allowlist_only', 'auto_enroll_review', 'auto_enroll_autopilot'].includes(mode)) throw new Error('不支持的个人消息接入模式');
- if (!['draft', 'send'].includes(welcomeSendMode)) throw new Error('不支持的欢迎语发送模式');
- if (welcomeEnabled && !welcomeText) throw new Error('启用欢迎语时必须填写欢迎内容');
- if (welcomeSendMode === 'send' && mode !== 'auto_enroll_autopilot') throw new Error('欢迎语直接发送仅可用于全自动新联系人接入');
- if (mode === 'auto_enroll_autopilot' && input.confirmation !== INTAKE_AUTOPILOT_CONFIRMATION) {
- throw new Error('开启全自动新联系人接入需要二次确认');
- }
- const confirmed = mode === 'auto_enroll_autopilot';
- target.db.setIntakePolicy({
- mode,
- welcomeEnabled,
- welcomeText,
- welcomeSendMode,
- autopilotConfirmation: confirmed ? INTAKE_AUTOPILOT_CONFIRMATION_VERSION : '',
- autopilotConfirmedAt: confirmed ? new Date().toISOString() : '',
- autopilotConfirmedBy: confirmed ? String(actor || 'human') : '',
- });
- target.db.audit({
- actor,
- action: 'personal_intake_policy_updated',
- detail: { mode, welcomeEnabled, welcomeSendMode, effectiveConversationMode: mode === 'auto_enroll_autopilot' ? 'autopilot' : 'review' },
- });
- return intakePolicyPayload(target);
- }
- function updateIntakePolicy(input = {}) {
- return {
- status: 'ok',
- assistantMessage: '个人消息接入策略已保存',
- data: updateIntakePolicyForWorkbench(workbench, input, 'human'),
- warnings: [],
- errors: [],
- };
- }
- async function retryOnboardingWelcome(identifier) {
- const value = String(identifier || '');
- const conversation = workbench.db.getConversation(value);
- const contactId = conversation?.contact_id || value;
- const result = await workbench.service.retryOnboardingWelcome(contactId, 'human');
- return {
- status: 'ok',
- assistantMessage: result.status === 'welcome_sent' ? '欢迎语重试已发送' : '欢迎语重试已处理',
- data: result,
- };
- }
- function displayableConversations() {
- return workbench.db.listConversations();
- }
- function getConversations() {
- const conversations = sortConversationsByRecency(displayableConversations().map(publicConversation));
- return {
- status: 'ok',
- data: {
- conversations,
- summary: {
- privateCount: conversations.filter(item => item.channelType === 'private').length,
- groupCount: conversations.filter(item => item.channelType === 'group').length,
- },
- },
- };
- }
- function getResponseMonitor() {
- const allowlist = new Set(workbench.config.qiwei.allowedSenders.map(String));
- const conversations = displayableConversations()
- .filter(row => allowlist.has(String(row.contact_id || '')))
- .flatMap(row => {
- const conversation = publicConversation(row);
- if (conversation.channelType === 'group') return [];
- return [{
- id: conversation.id,
- contactId: row.contact_id,
- displayName: conversation.displayName,
- messages: conversation.messages,
- handledMessageId: conversation.agentNotice?.entityId || '',
- }];
- });
- const account = activeAccountMetadata();
- const data = responseMonitor({
- projectRoot: PROJECT_ROOT,
- conversations,
- selfUserIds: [workbench.config.qiwei.selfUserId, account.userId],
- selfNames: [account.nickname],
- warningMinutes: workbench.config.qiwei.responseReminderMinutes,
- urgentMinutes: workbench.config.qiwei.responseUrgentMinutes,
- });
- return {
- status: 'ok',
- assistantMessage: data.summary.overdue
- ? `当前有 ${data.summary.overdue} 个会话超过回复时效,请客服优先处理。`
- : '当前监听范围内没有超时未回复会话。',
- summary: data.summary,
- data,
- warnings: [],
- errors: [],
- };
- }
- function updateCustomerProfile(conversationId, input = {}) {
- const conversation = workbench.db.getConversation(conversationId);
- if (!conversation) throw new Error('客户会话不存在');
- const current = workbench.db.getProfile(conversationId);
- const nextProfile = { ...(current.profile || {}) };
- const evidence = { ...(nextProfile.__evidence || {}) };
- const patch = input.profile && typeof input.profile === 'object' ? input.profile : {};
- const changedFields = [];
- for (const [field, rawValue] of Object.entries(patch)) {
- if (!field || field.startsWith('__')) continue;
- const value = typeof rawValue === 'string' ? rawValue.trim() : rawValue;
- if (value === '' || value === null || value === undefined) delete nextProfile[field];
- else nextProfile[field] = value;
- evidence[field] = {
- text: String(input.reason || '客户管理人工核对').trim(),
- sourceMessageId: null,
- source: 'human',
- updatedAt: new Date().toISOString(),
- };
- changedFields.push(field);
- }
- nextProfile.__evidence = evidence;
- const tags = input.tags === undefined
- ? current.tags
- : [...new Set((Array.isArray(input.tags) ? input.tags : String(input.tags || '').split(/[,,]/)).map(item => String(item).trim()).filter(Boolean))];
- const updated = workbench.db.updateProfile(conversationId, nextProfile, tags);
- const memory = workbench.service.memory.syncHumanProfile(conversationId, patch);
- workbench.db.audit({
- actor: 'human',
- action: 'customer_profile_updated',
- conversationId,
- entityId: conversationId,
- detail: { fields: changedFields, tagCount: tags.length, reason: String(input.reason || '').trim() },
- });
- const { __evidence, ...visibleProfile } = updated.profile || {};
- return {
- status: 'ok',
- assistantMessage: `客户“${conversation.contact_name || '未命名客户'}”的主档已更新。`,
- summary: { changedFields, tagCount: tags.length },
- data: { profile: visibleProfile, profileEvidence: __evidence || {}, tags, updatedAt: updated.updatedAt, memory },
- warnings: [],
- errors: [],
- };
- }
- function resolveConversationSyncScope(input = {}, allowlist = new Set(), db = workbench.db) {
- const conversationId = String(input.conversationId || '').trim();
- let contactId = String(input.contactId || '').trim();
- if (conversationId) {
- const conversation = db.getConversation(conversationId);
- if (!conversation) throw new Error('当前会话不存在,请刷新后重试');
- if (conversationChannelInfo(conversation, db.listMessages(conversation.id, 100)).channelType !== 'private') {
- throw new Error('当前采集仅支持个人聊天会话');
- }
- contactId = String(conversation.contact_id || '').trim();
- }
- if (contactId && !allowlist.has(contactId)) throw new Error('当前联系人不在个人消息白名单中');
- return {
- conversationId,
- contactId,
- contacts: contactId ? new Set([contactId]) : allowlist,
- scope: contactId ? 'conversation' : 'allowlist',
- };
- }
- async function syncConversations(input = {}) {
- const allowlist = new Set(workbench.config.qiwei.allowedSenders.map(String));
- if (!allowlist.size) throw new Error('测试联系人白名单为空,无法同步会话');
- const syncScope = resolveConversationSyncScope(input, allowlist);
- const account = await detectAccountStatus(true);
- if (!account.online) throw new Error('测试账号当前不在线,无法同步企微会话');
- const grouped = new Map();
- const seen = new Set();
- const seenSemantic = new Set();
- let cursor = 0;
- let pages = 0;
- let scannedMessages = 0;
- while (pages < workbench.config.qiwei.initialSyncMaxPages) {
- const result = await workbench.qiwei.syncMessages(cursor, workbench.config.qiwei.initialSyncLimit);
- const list = Array.isArray(result.syncMsgList) ? result.syncMsgList : [];
- scannedMessages += list.length;
- for (const message of list) {
- if (isGroupMessage(message)) continue;
- const senderId = String(message.senderId || '');
- const receiverId = String(message.receiverId || '');
- const contactId = syncScope.contacts.has(senderId) ? senderId : syncScope.contacts.has(receiverId) ? receiverId : '';
- const content = String(message.msgData?.content || '').trim();
- if (!contactId || !content || ![0, 1, 2].includes(Number(message.msgType))) continue;
- const inbound = senderId === contactId;
- const timestampSeconds = messageTimestamp(message.timestamp);
- if (!timestampSeconds) continue;
- const timestamp = new Date(timestampSeconds * 1000).toISOString();
- const externalId = String(message.msgServerId || message.msgUniqueIdentifier || '');
- const dedupeKey = externalId || `${contactId}|${inbound ? 'in' : 'out'}|${timestamp}|${content}`;
- if (seen.has(dedupeKey)) continue;
- seen.add(dedupeKey);
- const semanticKey = `${contactId}|${inbound ? 'in' : 'out'}|${timestamp}|${content}`;
- if (seenSemantic.has(semanticKey)) continue;
- seenSemantic.add(semanticKey);
- if (!grouped.has(contactId)) grouped.set(contactId, { contactName: '', messages: [] });
- const group = grouped.get(contactId);
- if (inbound && message.senderName) group.contactName = String(message.senderName);
- group.messages.push({
- externalId: externalId || null,
- inbound,
- content,
- timestamp,
- raw: { seq: message.seq, msgType: message.msgType, timestamp: message.timestamp, source: 'manual_sync' },
- });
- }
- const seqs = list.map(item => Number(item.seq)).filter(Number.isFinite);
- const next = Math.max(cursor, Number(result.travelSyncKey) || 0, seqs.length ? Math.max(...seqs) : 0);
- pages += 1;
- if (!list.length || next <= cursor) break;
- cursor = next;
- }
- let syncedMessages = 0;
- let removedImportedMessages = 0;
- for (const [contactId, group] of grouped.entries()) {
- const conversation = workbench.db.ensureConversation(contactId, group.contactName || '白名单企微联系人');
- removedImportedMessages += workbench.db.deleteImportedMessages(conversation.id, 'manual_sync');
- const ordered = group.messages.sort((a, b) => Date.parse(a.timestamp) - Date.parse(b.timestamp));
- const recent = ordered.slice(-20);
- for (const message of recent) {
- const inserted = workbench.db.insertMessage({
- conversationId: conversation.id,
- externalId: message.externalId,
- direction: message.inbound ? 'inbound' : 'outbound',
- senderType: message.inbound ? 'customer' : 'human',
- content: message.content,
- status: message.inbound ? 'received' : 'sent',
- createdAt: message.timestamp,
- raw: message.raw,
- });
- if (inserted.created) syncedMessages += 1;
- }
- }
- const duplicatesRemoved = workbench.db.cleanupInboundContentDuplicates(60);
- workbench.db.audit({
- actor: 'human',
- action: 'conversation_history_synced',
- detail: { scope: syncScope.scope, contactId: syncScope.contactId || null, conversations: grouped.size, insertedMessages: syncedMessages, removedImportedMessages, scannedMessages, duplicatesRemoved },
- });
- return {
- status: 'ok',
- assistantMessage: grouped.size
- ? syncScope.scope === 'conversation'
- ? '已采集当前个人会话的最近消息;只入库,不运行 Agent、不发送回复'
- : `已补采 ${grouped.size} 个白名单真实会话的最近消息;只入库,不运行 Agent、不发送回复`
- : syncScope.scope === 'conversation'
- ? '当前个人会话暂无可采集的文字消息,请先在企微中收发一条消息'
- : '未读取到白名单联系人的文字会话,请先在企微中与该联系人收发一条消息',
- data: {
- scope: syncScope.scope,
- conversationId: syncScope.conversationId || null,
- syncedConversationCount: grouped.size,
- syncedMessageCount: syncedMessages,
- scannedMessageCount: scannedMessages,
- conversations: sortConversationsByRecency(displayableConversations().map(publicConversation)),
- },
- };
- }
- function changeGlobalMode(mode, confirmation = '') {
- const selected = String(mode || '');
- requireAutopilotConfirmation(selected, confirmation);
- let update;
- if (['paused', 'monitor'].includes(selected)) update = { paused: true };
- else if (['review', 'suggest'].includes(selected)) update = { paused: false, defaultMode: 'review' };
- else if (selected === 'auto') update = { paused: false, defaultMode: 'auto' };
- else if (selected === 'autopilot') update = { paused: false, defaultMode: 'autopilot' };
- else if (selected === 'human') update = { paused: false, defaultMode: 'human' };
- else throw new Error('不支持的 Agent 模式');
- workbench.db.setSetting('listener_enabled', selected === 'paused' ? 'false' : 'true');
- const global = workbench.service.setGlobal(update);
- return {
- status: 'ok',
- assistantMessage: global.paused ? 'Agent 已全局暂停;仍可接收消息,但不会生成或发送回复' : `全局策略已切换为 ${global.defaultMode}`,
- data: { global, globalMode: global.paused ? 'paused' : global.defaultMode },
- };
- }
- function changeConversationMode(id, mode, confirmation = '') {
- const mapped = mode === 'agent' ? 'review' : mode === 'manual' ? 'human' : mode;
- requireAutopilotConfirmation(mapped, confirmation, 'conversation');
- const conversation = workbench.service.setConversationMode(id, mapped);
- const labels = { review: '待审核', auto: '高置信自动', autopilot: '全自动接管', human: '人工接管', paused: '会话暂停' };
- return { status: 'ok', assistantMessage: `会话已切换为${labels[mapped]}`, data: { conversation } };
- }
- async function approveReply(id, content) {
- const detail = workbench.service.conversationDetail(id);
- if (!detail) throw new Error('会话不存在');
- const pending = detail.drafts.find(item => item.status === 'pending');
- const result = pending
- ? await workbench.service.approveDraft(pending.id, { content, actor: 'human' })
- : { status: 'sent', message: await workbench.service.manualSend(id, content, 'human') };
- return { status: 'ok', assistantMessage: '回复已真实发送给白名单测试联系人', data: result };
- }
- async function approveDraft(draftId, content) {
- const result = await workbench.service.approveDraft(draftId, { content, actor: 'human' });
- return { status: 'ok', assistantMessage: '审核通过,回复已真实发送且只发送一次', data: result };
- }
- function rejectDraft(draftId, reason) {
- const result = workbench.service.rejectDraft(draftId, { reason, actor: 'human' });
- return { status: 'ok', assistantMessage: '草稿已驳回,不会发送给客户', data: result };
- }
- async function regenerateDraft(draftId) {
- const result = await workbench.service.regenerateDraft(draftId, 'human');
- return {
- status: 'ok',
- assistantMessage: result.status === 'pending_review' ? 'Agent 已重新生成待审核草稿' : (result.error || 'Agent 重新生成未完成'),
- data: result,
- };
- }
- async function generateLatestDraft(conversationId) {
- const result = await workbench.service.generateLatestDraft(conversationId, 'human');
- return {
- status: 'ok',
- assistantMessage: result.assistantMessage || (result.status === 'pending_review' ? 'Agent 已生成待审核草稿' : (result.error || 'Agent 未生成草稿')),
- data: result,
- };
- }
- async function manualSend(conversationId, content) {
- const message = await workbench.service.manualSend(conversationId, content, 'human');
- return { status: 'ok', assistantMessage: '人工回复已真实发送给白名单测试联系人', data: { message } };
- }
- function getVoiceStatus() {
- return { status: 'ok', data: workbench.voice.status() };
- }
- async function enrollVoice(input = {}) {
- const result = await workbench.voice.enroll({
- filePath: input.filePath,
- originalName: input.originalName,
- mime: input.mime,
- });
- workbench.db.audit({
- actor: 'human',
- action: 'voice_profile_enrolled',
- detail: { duration: result.profile?.duration, sampleRate: result.profile?.sampleRate },
- });
- return {
- status: 'ok',
- assistantMessage: '本人声音已初始化,可以生成企微语音',
- data: workbench.voice.status(),
- };
- }
- function revokeVoiceProfile() {
- const result = workbench.voice.revoke();
- workbench.db.audit({ actor: 'human', action: 'voice_profile_revoked', detail: {} });
- return { status: 'ok', assistantMessage: '声音档案已删除', data: result };
- }
- function voiceConversationContext(conversationId) {
- const conversation = workbench.db.getConversation(conversationId);
- if (!conversation) throw new Error('会话不存在');
- workbench.service.requireAllowed(conversation.contact_id);
- workbench.service.requirePrivateConversation(conversation.id);
- const inbound = workbench.db.getLatestInbound(conversationId);
- return { conversation, context: inbound?.content || '' };
- }
- function pendingVoiceDraft(db, conversationId, draftId = '') {
- const selectedId = String(draftId || '').trim();
- const draft = selectedId
- ? db.getDraft(selectedId)
- : db.listDrafts({ conversationId, status: 'pending', limit: 1 })[0];
- if (!draft && selectedId) throw new Error('待审核草稿不存在或已被处理,请刷新后重试');
- if (!draft) return null;
- if (draft.conversation_id !== conversationId) throw new Error('待审核草稿与当前会话不匹配');
- if (draft.status !== 'pending') throw new Error('待审核草稿已被处理,请刷新后重试');
- return draft;
- }
- function markVoiceDraftSent(db, draft, { content, messageId, actor = 'human:voice' }) {
- if (!draft) return null;
- return db.updateDraft(draft.id, {
- content,
- status: 'sent',
- reviewed_at: new Date().toISOString(),
- reviewer: actor,
- sent_message_id: messageId,
- error: null,
- });
- }
- async function sendClonedVoice(conversationId, input = {}) {
- const { conversation, context } = voiceConversationContext(conversationId);
- const content = String(input.content || '').trim();
- const draft = pendingVoiceDraft(workbench.db, conversationId, input.draftId);
- const result = await workbench.voice.synthesizeAndSend({
- text: content,
- context,
- tone: input.tone || 'auto',
- toId: conversation.contact_id,
- confirmed: input.confirmed,
- });
- if (result.sendResult?.isSendSuccess === false || Number(result.sendResult?.isSendSuccess) === 0) {
- throw new Error('企微未确认语音发送成功');
- }
- const message = workbench.db.insertMessage({
- conversationId,
- direction: 'outbound',
- senderType: 'human',
- content: `[语音 · ${result.tone.label}] ${content}`,
- status: 'sent',
- raw: {
- msgType: 16,
- tone: result.tone,
- duration: result.duration,
- fileSize: result.fileSize,
- audioPath: result.audioPath,
- sendResult: result.sendResult,
- },
- }).message;
- const resolvedDraft = markVoiceDraftSent(workbench.db, draft, { content, messageId: message.id });
- workbench.db.audit({
- actor: 'human',
- action: 'cloned_voice_sent',
- conversationId,
- entityId: message.id,
- detail: {
- tone: result.tone.id,
- automatic: result.tone.automatic,
- duration: result.duration,
- msgServerId: result.sendResult?.msgServerId,
- draftId: resolvedDraft?.id || null,
- },
- });
- const publicMessage = { ...message };
- delete publicMessage.raw_json;
- return {
- status: 'ok',
- assistantMessage: `已用${result.tone.label}语气发送企微语音`,
- data: { message: publicMessage, draft: resolvedDraft, tone: result.tone, duration: result.duration, sendResult: result.sendResult },
- };
- }
- function sentVoiceRuns(voiceRoot = categoryDir('voice')) {
- const runs = [];
- if (!fs.existsSync(voiceRoot)) return runs;
- for (const dateEntry of fs.readdirSync(voiceRoot, { withFileTypes: true })) {
- if (!dateEntry.isDirectory() || !/^\d{4}-\d{2}-\d{2}$/.test(dateEntry.name)) continue;
- const dateDir = path.join(voiceRoot, dateEntry.name);
- for (const runEntry of fs.readdirSync(dateDir, { withFileTypes: true })) {
- if (!runEntry.isDirectory() || !runEntry.name.includes('-clone')) continue;
- const runDir = path.join(dateDir, runEntry.name);
- const manifestPath = path.join(runDir, 'manifest.json');
- const audioPath = path.join(runDir, 'speech.wav');
- if (!fs.existsSync(manifestPath) || !fs.existsSync(audioPath)) continue;
- const manifest = parseJson(fs.readFileSync(manifestPath, 'utf8'), {});
- if (manifest.type !== 'voice-clone' || manifest.outcome === 'failed') continue;
- const createdAt = Date.parse(manifest.createdAt);
- if (!Number.isFinite(createdAt)) continue;
- const textHash = String(manifest.textSha256 || (manifest.text
- ? crypto.createHash('sha256').update(String(manifest.text)).digest('hex')
- : ''));
- if (!textHash) continue;
- runs.push({
- audioPath,
- createdAt,
- duration: Number(manifest.audio?.duration || manifest.duration) || 0,
- textHash,
- });
- }
- }
- return runs;
- }
- function backfillSentVoiceAudioPaths(db, options = {}) {
- const available = sentVoiceRuns(options.voiceRoot).map(item => ({ ...item, used: false }));
- if (!available.length) return 0;
- let updated = 0;
- for (const conversation of db.listConversations()) {
- for (const message of db.listMessages(conversation.id, 1000)) {
- const raw = parseJson(message.raw_json, {});
- if (message.direction !== 'outbound' || message.status !== 'sent'
- || Number(raw.msgType) !== 16 || raw.audioPath) continue;
- const spokenText = String(message.content || '').replace(/^\[语音[^\]]*\]\s*/, '');
- const textHash = crypto.createHash('sha256').update(spokenText).digest('hex');
- const sentAt = Date.parse(message.created_at);
- const match = available
- .filter(item => !item.used && item.textHash === textHash
- && Math.abs(item.createdAt - sentAt) <= 10 * 60 * 1000
- && (!item.duration || !raw.duration || Math.abs(item.duration - Number(raw.duration)) <= 1.5))
- .sort((a, b) => Math.abs(a.createdAt - sentAt) - Math.abs(b.createdAt - sentAt))[0];
- if (!match) continue;
- match.used = true;
- db.updateMessageRaw(message.id, { ...raw, audioPath: match.audioPath });
- updated += 1;
- }
- }
- return updated;
- }
- function getSentVoiceAudio(messageId) {
- const message = workbench.db.getMessage(String(messageId || ''));
- const raw = parseJson(message?.raw_json, {});
- if (!message || message.direction !== 'outbound' || message.status !== 'sent'
- || Number(raw.msgType) !== 16 || !raw.audioPath) {
- throw new Error('已发送语音不存在或不可播放');
- }
- return { filePath: raw.audioPath, duration: Number(raw.duration) || 0 };
- }
- async function updateCustomerTask(taskId, input = {}) {
- if (!['open', 'in_progress', 'done', 'dismissed'].includes(String(input.status || ''))) throw new Error('不支持的客户待办状态');
- const existing = workbench.db.getCustomerTask(taskId);
- if (!existing) throw new Error('客户待办不存在');
- let officialResult = null;
- if (input.status === 'done' && existing.official_todo_id) {
- officialResult = await completeTodoKnowledge(existing.official_todo_id);
- }
- const officialOk = !officialResult || officialResult.status === 'ok';
- const task = workbench.db.updateCustomerTask(taskId, {
- status: input.status,
- resolution_reason: input.status === 'done' ? 'human_completed' : input.status === 'dismissed' ? 'human_dismissed' : '',
- ...(existing.official_todo_id ? { official_sync_status: officialOk ? input.status : 'error' } : {}),
- });
- if (!task) throw new Error('客户待办不存在');
- workbench.db.audit({ actor: 'human', action: 'customer_task_updated', conversationId: task.conversation_id, entityId: task.id, detail: { status: task.status, officialSynced: Boolean(existing.official_todo_id), officialOk } });
- return {
- status: 'ok',
- assistantMessage: existing.official_todo_id && !officialOk
- ? `本地待办已更新为 ${task.status},但企微官方待办同步失败,请稍后重试`
- : `客户待办已更新为 ${task.status}${existing.official_todo_id ? ',企微官方待办已同步' : ''}`,
- data: { task, officialResult },
- warnings: existing.official_todo_id && !officialOk ? ['企微官方待办状态尚未同步'] : [],
- };
- }
- async function syncCustomerTaskToOfficialTodo(taskId, input = {}) {
- const syncCurrentAccountTask = createCustomerTaskOfficialSync({
- db: workbench.db,
- searchTodoUsers,
- createTodoKnowledge,
- });
- return syncCurrentAccountTask(taskId, input);
- }
- function updateCustomerAlert(alertId, input = {}) {
- if (!['open', 'acknowledged', 'resolved', 'dismissed'].includes(String(input.status || ''))) throw new Error('不支持的客户预警状态');
- const alert = workbench.db.updateCustomerAlert(alertId, { status: input.status });
- if (!alert) throw new Error('客户预警不存在');
- workbench.db.audit({ actor: 'human', action: 'customer_alert_updated', conversationId: alert.conversation_id, entityId: alert.id, detail: { status: alert.status } });
- return { status: 'ok', assistantMessage: `客户预警已更新为 ${alert.status}`, data: { alert } };
- }
- function addCustomerMemory(conversationId, input = {}) {
- const result = workbench.service.addCustomerMemory(conversationId, input, 'human');
- return { status: 'ok', assistantMessage: '客户长期记忆已添加,并将在后续回复中生效', data: result };
- }
- function updateCustomerMemory(memoryId, input = {}) {
- const action = String(input.action || 'update');
- if (action === 'forget') {
- const result = workbench.service.forgetCustomerMemory(memoryId, 'human');
- return { status: 'ok', assistantMessage: '客户记忆已彻底遗忘', data: result };
- }
- const patch = action === 'confirm'
- ? { type: 'fact', status: 'active', confidence: 1 }
- : action === 'reject'
- ? { status: 'rejected' }
- : {
- content: input.content,
- type: input.type,
- status: input.status,
- importance: input.importance,
- expiresAt: input.expiresAt,
- };
- const sanitized = Object.fromEntries(Object.entries(patch).filter(([, value]) => value !== undefined));
- const result = workbench.service.updateCustomerMemory(memoryId, sanitized, 'human');
- const message = action === 'confirm' ? '推断已人工确认为客户事实' : action === 'reject' ? '客户记忆已拒绝,不再进入上下文' : '客户记忆已更新';
- return { status: 'ok', assistantMessage: message, data: result };
- }
- function getAudit(limit = 200) {
- return { status: 'ok', data: { audit: workbench.db.listAudit(Math.max(1, Math.min(500, Number(limit) || 200))) } };
- }
- async function startListenerForWorkbench(target, account, options = {}) {
- if (!account.online) throw new Error(`${account.nickname || '当前账号'}不在线,无法启动真实消息监听`);
- refreshAllowedSendersFromEnv(target.config.qiwei, options.envFile || ENV_FILE);
- if (options.automatic === true && target.db?.getSetting?.('listener_enabled', 'true') === 'false') {
- return { status: 'ok', assistantMessage: 'AI 监听保持人工关闭状态', data: { running: false, disabled: true } };
- }
- target.db?.setSetting?.('listener_enabled', 'true');
- target.service?.setGlobal?.({ paused: false }, options.automatic ? 'runtime:auto-start' : 'human');
- const status = await target.poller.start();
- return {
- status: 'ok',
- assistantMessage: 'AI 监听已启动:白名单私聊按当前策略处理;已确认客户群自动生成待审核草稿,不会自动群发',
- data: status,
- };
- }
- async function ingestRuntimeMessage(message = {}, options = {}) {
- return ingestMessageForWorkbench({
- config: workbench.config.qiwei,
- db: workbench.db,
- service: workbench.service,
- qiwei: workbench.qiwei,
- }, message, String(options.source || 'runtime'));
- }
- async function startListener(options = {}) {
- const product = getProductMode();
- if (product.mode === 'enterprise') {
- const runtime = listenerStateWithRuntime(product, workbench.poller.status());
- if (!runtime.running) throw new Error('Enterprise Relay runtime is not running. Start the Qiwei runtime first.');
- workbench.service.setGlobal({ paused: false, defaultMode: 'review' });
- return {
- status: 'ok',
- assistantMessage: 'Enterprise Relay keeps collecting messages; AI review and reply processing is enabled.',
- data: runtime,
- };
- }
- const account = await detectAccountStatus(true);
- return startListenerForWorkbench(workbench, account, options);
- }
- function applyManualTakeover(target) {
- target.service.setGlobal({ paused: true, defaultMode: 'review' });
- for (const conversation of target.db.listConversations()) {
- target.service.setConversationMode(conversation.id, 'human');
- }
- }
- function stopListener({ preserveAgentState = false } = {}) {
- const product = getProductMode();
- if (product.mode === 'enterprise') {
- if (!preserveAgentState) {
- workbench.service.setGlobal({ paused: true, defaultMode: 'review' });
- for (const conversation of workbench.db.listConversations()) {
- workbench.service.setConversationMode(conversation.id, 'human');
- }
- }
- return {
- status: 'ok',
- assistantMessage: preserveAgentState
- ? 'Enterprise Relay runtime stopped; Agent modes were preserved.'
- : 'Enterprise Relay continues collecting messages; AI replies are paused and conversations are in human mode.',
- data: { running: false, relayRunning: true, transport: 'server_relay' },
- };
- }
- const status = workbench.poller.stop();
- if (!preserveAgentState) {
- workbench.db.setSetting('listener_enabled', 'false');
- applyManualTakeover(workbench);
- }
- return {
- status: 'ok',
- assistantMessage: preserveAgentState
- ? 'Qiwei polling runtime stopped; Agent modes were preserved.'
- : 'AI 监听已关闭,现有会话已切换为人工接管',
- data: status,
- };
- }
- function getAgentRuntimeConfig() {
- return { ...workbench.config.agent };
- }
- module.exports = {
- switchActiveAccount,
- getAgentStatus,
- getAllowlistCandidates,
- updateAllowlist,
- addAllowlistContacts,
- getIntakePolicy,
- updateIntakePolicy,
- retryOnboardingWelcome,
- getConversations,
- getResponseMonitor,
- updateCustomerProfile,
- syncConversations,
- changeGlobalMode,
- changeConversationMode,
- approveReply,
- approveDraft,
- rejectDraft,
- regenerateDraft,
- generateLatestDraft,
- manualSend,
- getVoiceStatus,
- enrollVoice,
- revokeVoiceProfile,
- sendClonedVoice,
- getSentVoiceAudio,
- updateCustomerTask,
- syncCustomerTaskToOfficialTodo,
- updateCustomerAlert,
- addCustomerMemory,
- updateCustomerMemory,
- getAudit,
- ingestRuntimeMessage,
- ingestWebhookMessage,
- startListener,
- stopListener,
- getAgentRuntimeConfig,
- createWorkbench,
- __testing: { loadAgentConfig, normalizeAllowlistIds, normalizeAllowlistContact, refreshAllowedSendersFromEnv, hydrateAccountAllowlist, updateAllowlistForWorkbench, resolveConversationSyncScope, applyManualTakeover, requireAutopilotConfirmation, accountRuntimeKey, accountWorkbenchOverrides, activeAccountMetadata, FmodeQiweiClient, QiweiAgentPoller, outboundContactId, recordOutboundMessage, syncOutboundNotification, ingestMessageForWorkbench, conversationChannelInfo, publicConversation, backfillCustomerIntelligence, backfillCustomerMemory, backfillSentVoiceAudioPaths, sentVoiceRuns, pendingVoiceDraft, markVoiceDraftSent, startListenerForWorkbench, intakePolicyPayload, updateIntakePolicyForWorkbench },
- };
|