agent-service.js 85 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972
  1. const fs = require('fs');
  2. const path = require('path');
  3. const crypto = require('crypto');
  4. const { PACKAGE_ROOT, WORKSPACE_ROOT, latestPath, categoryDir, outputsRoot } = require('../core/output-paths');
  5. const { AgentWorkbenchDb } = require('../core/agent-workbench-db');
  6. const { AgentKnowledgeStore } = require('../core/agent-knowledge');
  7. const { QiweiAgentRuntime, extractExplicitCustomerIntelligence } = require('../core/agent-runtime');
  8. const { getCustomerSessionGuide } = require('../core/agent-session-guide');
  9. const { AgentWorkbenchService } = require('../core/agent-workbench-service');
  10. const { friendlyAgentError } = require('../core/agent-error-message');
  11. const { VoiceCloneService } = require('../core/voice-clone-service');
  12. const {
  13. searchTodoUsers,
  14. createTodoKnowledge,
  15. completeTodoKnowledge,
  16. } = require('./official-office-knowledge-service');
  17. const { createCustomerTaskOfficialSync } = require('../core/customer-task-official-sync');
  18. const { messageTimestamp, roomIdOf, isGroupMessage, messageContent, evaluatePolledMessage, normalizePersonalIntakeMode } = require('../core/agent-poller-policy');
  19. const { sortConversationsByRecency, maxConversationTimestamp } = require('../core/conversation-order');
  20. const { saveQiweiClientConfig, setActiveQiweiContext, readFmodeVoiceToken } = require('../core/credentials');
  21. const { FmodeQiweiClient } = require('../providers/fmode-agent-transport');
  22. const { responseMonitor } = require('./response-monitor-service');
  23. const { normalizeAllowlistIds, normalizeAllowlistContact } = require('../core/allowlist-config');
  24. const { getProductMode } = require('../core/product-mode');
  25. const { getRelayDaemonStatus } = require('../core/relay-daemon');
  26. const PROJECT_ROOT = WORKSPACE_ROOT;
  27. const ENV_FILE = path.join(PROJECT_ROOT, '.env.local');
  28. const INTAKE_AUTOPILOT_CONFIRMATION = 'ENABLE_AUTO_ENROLL_AUTOPILOT';
  29. const INTAKE_AUTOPILOT_CONFIRMATION_VERSION = 'v1';
  30. const AUTOPILOT_CONFIRMATION = 'ENABLE_AUTOPILOT';
  31. function requireAutopilotConfirmation(mode, confirmation, scope = 'global') {
  32. if (mode !== 'autopilot' || confirmation === AUTOPILOT_CONFIRMATION) return;
  33. throw new Error(scope === 'conversation' ? '开启会话全自动接管需要二次确认' : '开启全自动接管需要二次确认');
  34. }
  35. function readEnvFile(filePath) {
  36. try {
  37. const env = {};
  38. for (const rawLine of fs.readFileSync(filePath, 'utf8').replace(/^\uFEFF/, '').split(/\r?\n/)) {
  39. const line = rawLine.trim();
  40. if (!line || line.startsWith('#') || !line.includes('=')) continue;
  41. const index = line.indexOf('=');
  42. const key = line.slice(0, index).trim();
  43. let value = line.slice(index + 1).trim();
  44. if ((value.startsWith('"') && value.endsWith('"')) || (value.startsWith("'") && value.endsWith("'"))) value = value.slice(1, -1);
  45. env[key] = value;
  46. }
  47. return env;
  48. } catch {
  49. return {};
  50. }
  51. }
  52. function refreshAllowedSendersFromEnv(config = {}, envFile = ENV_FILE) {
  53. if (config.usePersistedAllowlist === true) {
  54. const combined = normalizeAllowlistIds([...(config.manualAllowedSenders || []), ...(config.autoEnrolledSenders || [])]);
  55. const current = Array.isArray(config.allowedSenders) ? config.allowedSenders.map(String) : [];
  56. const changed = combined.length !== current.length || combined.some((id, index) => id !== current[index]);
  57. if (changed) config.allowedSenders = combined;
  58. return { changed, count: combined.length };
  59. }
  60. const env = readEnvFile(envFile);
  61. if (!Object.prototype.hasOwnProperty.call(env, 'QIWEI_AUTO_REPLY_ALLOWED_SENDERS')) {
  62. return { changed: false, count: Array.isArray(config.allowedSenders) ? config.allowedSenders.length : 0 };
  63. }
  64. const next = normalizeAllowlistIds(env.QIWEI_AUTO_REPLY_ALLOWED_SENDERS);
  65. const combined = normalizeAllowlistIds([...next, ...(config.autoEnrolledSenders || [])]);
  66. const current = Array.isArray(config.allowedSenders) ? config.allowedSenders.map(String) : [];
  67. const changed = combined.length !== current.length || combined.some((id, index) => id !== current[index]);
  68. config.manualAllowedSenders = next;
  69. if (changed) config.allowedSenders = combined;
  70. return { changed, count: combined.length };
  71. }
  72. function storedAllowlistIds(db, key, fallback = []) {
  73. const raw = db.getSetting(key, '');
  74. if (!raw) {
  75. const ids = normalizeAllowlistIds(fallback);
  76. db.setSetting(key, JSON.stringify(ids));
  77. return ids;
  78. }
  79. try { return normalizeAllowlistIds(JSON.parse(raw)); } catch { return []; }
  80. }
  81. function hydrateAccountAllowlist(config, db) {
  82. const manual = storedAllowlistIds(db, 'manual_allowed_sender_ids', config.qiwei.manualAllowedSenders || config.qiwei.allowedSenders || []);
  83. const automatic = db.getAutoEnrolledContactIds();
  84. config.qiwei.manualAllowedSenders = manual;
  85. config.qiwei.autoEnrolledSenders = automatic;
  86. config.qiwei.allowedSenders = normalizeAllowlistIds([...manual, ...automatic]);
  87. config.qiwei.usePersistedAllowlist = true;
  88. return config.qiwei.allowedSenders;
  89. }
  90. function readRuntimeState() {
  91. try {
  92. const filePath = path.join(outputsRoot(), 'runtime', 'qiwei-runtime.json');
  93. return JSON.parse(fs.readFileSync(filePath, 'utf8').replace(/^\uFEFF/, ''));
  94. } catch {
  95. return {};
  96. }
  97. }
  98. function runtimeProcessAlive(pid) {
  99. const numericPid = Number(pid);
  100. if (!Number.isInteger(numericPid) || numericPid <= 0) return false;
  101. try {
  102. process.kill(numericPid, 0);
  103. return true;
  104. } catch {
  105. return false;
  106. }
  107. }
  108. function listenerStateWithRuntime(product, localState) {
  109. const runtime = readRuntimeState();
  110. if (!runtimeProcessAlive(runtime.pid) || !['starting', 'running'].includes(runtime.status)) return localState;
  111. const component = product.mode === 'enterprise'
  112. ? runtime.components?.enterpriseRelay
  113. : runtime.components?.personalPolling;
  114. if (!component) return localState;
  115. return {
  116. ...localState,
  117. running: Boolean(localState.running || ['starting', 'running'].includes(component.status)),
  118. lastError: component.lastError || localState.lastError || '',
  119. transport: runtime.transport,
  120. runtimePid: runtime.pid,
  121. runtimeStatus: component.status,
  122. };
  123. }
  124. function readClaudeSettingsEnv() {
  125. const result = {};
  126. const home = process.env.USERPROFILE || process.env.HOME || '';
  127. for (const filePath of [path.join(home, '.claude', 'settings.json'), path.join(home, '.claude', 'settings.local.json')]) {
  128. try {
  129. const parsed = JSON.parse(fs.readFileSync(filePath, 'utf8').replace(/^\uFEFF/, ''));
  130. for (const [key, item] of Object.entries(parsed.env || {})) {
  131. if (!result[key] && typeof item === 'string' && item.trim()) result[key] = item.trim();
  132. }
  133. } catch {}
  134. }
  135. return result;
  136. }
  137. const fileEnv = readEnvFile(ENV_FILE);
  138. const claudeEnv = readClaudeSettingsEnv();
  139. function value(name, fallback = '') {
  140. const candidates = [process.env[name], fileEnv[name], claudeEnv[name], fallback];
  141. return String(candidates.find(item => typeof item === 'string' && item.trim()) ?? '').trim();
  142. }
  143. function bool(name, fallback = false) {
  144. return /^(1|true|yes|on)$/i.test(value(name, fallback ? 'true' : 'false'));
  145. }
  146. function number(name, fallback, min = -Infinity, max = Infinity) {
  147. const parsed = Number(value(name, String(fallback)));
  148. return Math.min(max, Math.max(min, Number.isFinite(parsed) ? parsed : fallback));
  149. }
  150. function resolvePath(input, fallback) {
  151. const selected = input || fallback;
  152. return selected ? path.resolve(PROJECT_ROOT, selected) : '';
  153. }
  154. function loadAgentConfig(overrides = {}) {
  155. const provider = value('AGENT_PROVIDER', 'claude-code');
  156. const anthropic = provider === 'anthropic';
  157. const claudeCode = provider === 'claude-code';
  158. const manualAllowedSenders = normalizeAllowlistIds(value('QIWEI_AUTO_REPLY_ALLOWED_SENDERS'));
  159. const baseConfig = {
  160. accountKey: accountRuntimeKey({ uid: value('QIWEI_UID') || value('QIWE_UID'), guid: value('QIWEI_GUID') || value('QIWE_GUID') }),
  161. dbPath: resolvePath(value('QIWEI_AGENT_DB_PATH'), latestPath('messages', 'agent-workbench.db')),
  162. legacyDbPath: path.resolve(PACKAGE_ROOT, '..', '..', 'qiwei-agent-workbench', 'data', 'workbench.db'),
  163. globalDefaultPaused: bool('QIWEI_AGENT_GLOBAL_DEFAULT_PAUSED', false),
  164. conversationDefaultMode: value('QIWEI_AGENT_DEFAULT_MODE', 'review'),
  165. autoSendConfidence: number('QIWEI_AGENT_AUTO_SEND_CONFIDENCE', 0.88, 0, 1),
  166. intake: {
  167. mode: normalizePersonalIntakeMode(value('QIWEI_PERSONAL_INTAKE_MODE', 'allowlist_only')),
  168. welcomeEnabled: bool('QIWEI_WELCOME_ENABLED', false),
  169. welcomeText: value('QIWEI_WELCOME_TEXT', '您好,已经收到您的消息,我会尽快为您处理。'),
  170. welcomeSendMode: value('QIWEI_WELCOME_SEND_MODE', 'draft') === 'send' ? 'send' : 'draft',
  171. effectiveAutopilotMode: 'autopilot',
  172. },
  173. knowledgeDir: resolvePath(
  174. value('QIWEI_AGENT_KNOWLEDGE_DIR'),
  175. fs.existsSync(path.join(PROJECT_ROOT, 'knowledge')) ? path.join(PROJECT_ROOT, 'knowledge') : path.join(PACKAGE_ROOT, 'knowledge')
  176. ),
  177. contextFiles: value('QIWEI_AGENT_CONTEXT_FILES', 'personality.md,context.md,rules.md').split(',').map(item => item.trim()).filter(Boolean),
  178. contextCharLimit: number('QIWEI_AGENT_CONTEXT_CHAR_LIMIT', 6000, 1000, 20000),
  179. memory: {
  180. enabled: bool('QIWEI_AGENT_MEMORY_ENABLED', true),
  181. coreCharLimit: number('QIWEI_AGENT_MEMORY_CORE_CHAR_LIMIT', 4000, 500, 8000),
  182. recallLimit: number('QIWEI_AGENT_MEMORY_RECALL_LIMIT', 6, 0, 20),
  183. recentMessageLimit: number('QIWEI_AGENT_MEMORY_RECENT_MESSAGES', 12, 4, 40),
  184. historyScanLimit: number('QIWEI_AGENT_MEMORY_HISTORY_SCAN_LIMIT', 240, 20, 1000),
  185. extractionIntervalMs: number('QIWEI_AGENT_MEMORY_EXTRACTION_INTERVAL_MS', 1000, 250, 60000),
  186. extractionRetryBaseMs: number('QIWEI_AGENT_MEMORY_EXTRACTION_RETRY_BASE_MS', 1000, 250, 60000),
  187. extractionMaxAttempts: number('QIWEI_AGENT_MEMORY_EXTRACTION_MAX_ATTEMPTS', 3, 1, 10),
  188. },
  189. agent: {
  190. provider,
  191. apiKey: value('AGENT_API_KEY') || value(anthropic ? 'ANTHROPIC_AUTH_TOKEN' : 'OPENAI_API_KEY') || value('QIWEI_AUTO_REPLY_AI_KEY'),
  192. 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(/\/$/, ''),
  193. model: value('AGENT_MODEL') || value(anthropic || claudeCode ? 'ANTHROPIC_MODEL' : 'OPENAI_MODEL') || value('QIWEI_AUTO_REPLY_AI_MODEL') || (anthropic || claudeCode ? 'sonnet' : 'gpt-4.1-mini'),
  194. maxToolRounds: number('AGENT_MAX_TOOL_ROUNDS', 4, 1, 8),
  195. promptCharLimit: number('QIWEI_AGENT_PROMPT_CHAR_LIMIT', 16000, 2000, 50000),
  196. systemPromptCharLimit: number('QIWEI_AGENT_SYSTEM_PROMPT_CHAR_LIMIT', 12000, 2000, 50000),
  197. qualityPassScore: number('QIWEI_AGENT_QUALITY_PASS_SCORE', 82, 60, 100),
  198. businessGoal: value('QIWEI_AGENT_BUSINESS_GOAL'),
  199. stagePlaybook: value('QIWEI_AGENT_STAGE_PLAYBOOK'),
  200. claudeExecutable: value('CLAUDE_CODE_EXECUTABLE') || path.join(path.dirname(process.execPath), 'node_modules', '@anthropic-ai', 'claude-code', 'bin', 'claude.exe'),
  201. claudeWorkdir: resolvePath(value('CLAUDE_CODE_WORKDIR'), PROJECT_ROOT),
  202. claudeSessionFile: resolvePath(value('CLAUDE_CODE_SESSION_FILE'), latestPath('messages', 'claude-code-sessions.json')),
  203. claudeProjectId: value('QIWEI_AGENT_PROJECT_ID') || crypto.createHash('sha256').update(PROJECT_ROOT).digest('hex').slice(0, 16),
  204. claudeMainSessionId: value('QIWEI_AGENT_MAIN_SESSION_ID'),
  205. claudeTimeoutMs: number('CLAUDE_CODE_TIMEOUT_MS', 120000, 15000, 300000),
  206. claudeMaxBudgetUsd: number('CLAUDE_CODE_MAX_BUDGET_USD', 1, 0.05, 5),
  207. claudeRetryMaxBudgetUsd: number('CLAUDE_CODE_RETRY_MAX_BUDGET_USD', 3, 0.1, 5),
  208. claudeBare: bool('CLAUDE_CODE_BARE', true),
  209. claudeEffort: value('CLAUDE_CODE_EFFORT', 'low'),
  210. claudeTools: value('CLAUDE_CODE_ALLOWED_TOOLS', 'Read,Glob,Grep'),
  211. claudeSessionMaxTurns: number('CLAUDE_CODE_SESSION_MAX_TURNS', 12, 1, 100),
  212. claudeSessionMaxAgeMs: number('CLAUDE_CODE_SESSION_MAX_AGE_MS', 86400000, 60000, 604800000),
  213. },
  214. qiwei: {
  215. transport: 'fmode-gateway',
  216. authToken: value('QIWEI_AUTH_TOKEN'),
  217. uid: value('QIWEI_UID') || value('QIWE_UID'),
  218. guid: value('QIWEI_GUID') || value('QIWE_GUID'),
  219. apiBase: value('QIWEI_API_BASE') || value('QIWE_API_BASE'),
  220. userId: '',
  221. nickname: '',
  222. corpName: '',
  223. allowedSenders: [...manualAllowedSenders],
  224. manualAllowedSenders: [...manualAllowedSenders],
  225. autoEnrolledSenders: [],
  226. accountKey: '',
  227. intakeAutopilotConversationMode: 'autopilot',
  228. selfUserId: value('QIWEI_AUTO_REPLY_SELF_USER_ID'),
  229. intervalMs: number('QIWEI_AUTO_REPLY_INTERVAL_MS', 10000, 3000, 60000),
  230. initialSyncLimit: number('QIWEI_AGENT_INITIAL_SYNC_LIMIT', 5000, 100, 5000),
  231. initialSyncMaxPages: number('QIWEI_AGENT_INITIAL_SYNC_MAX_PAGES', 200, 10, 500),
  232. startupGraceSeconds: number('QIWEI_AGENT_STARTUP_GRACE_SECONDS', 10, 0, 60),
  233. trustConfiguredOnStatusError: bool('QIWEI_TRUST_CONFIGURED_ON_STATUS_ERROR', false),
  234. responseReminderMinutes: number('QIWEI_RESPONSE_REMINDER_MINUTES', 15, 1, 1440),
  235. responseUrgentMinutes: number('QIWEI_RESPONSE_URGENT_MINUTES', 60, 1, 10080),
  236. },
  237. voice: {
  238. endpoint: value('QIWEI_VOICE_ENDPOINT', 'https://server.fmode.cn/api/voice/indextts2'),
  239. authToken: readFmodeVoiceToken({
  240. voiceAuthToken: value('QIWEI_VOICE_AUTH_TOKEN'),
  241. endpoint: value('QIWEI_VOICE_ENDPOINT', 'https://server.fmode.cn/api/voice/indextts2'),
  242. }),
  243. model: 'fmode-voice',
  244. requestTimeoutMs: number('QIWEI_TTS_TIMEOUT_MS', 180000, 15000, 300000),
  245. },
  246. };
  247. const config = {
  248. ...baseConfig,
  249. ...overrides,
  250. memory: { ...baseConfig.memory, ...(overrides.memory || {}) },
  251. intake: { ...baseConfig.intake, ...(overrides.intake || {}) },
  252. agent: { ...baseConfig.agent, ...(overrides.agent || {}) },
  253. qiwei: { ...baseConfig.qiwei, ...(overrides.qiwei || {}) },
  254. voice: { ...baseConfig.voice, ...(overrides.voice || {}) },
  255. };
  256. config.qiwei.accountKey = config.accountKey;
  257. config.qiwei.intakeAutopilotConversationMode = config.intake.effectiveAutopilotMode;
  258. config.agent.claudeAddDirs = [config.knowledgeDir].filter(Boolean);
  259. return config;
  260. }
  261. function accountRuntimeKey(input = {}) {
  262. const source = String(input.uid || input.guid || input.userId || 'default').trim();
  263. return crypto.createHash('sha256').update(source || 'default').digest('hex').slice(0, 16);
  264. }
  265. function accountWorkbenchOverrides(input = {}) {
  266. const storageKey = accountRuntimeKey(input);
  267. return {
  268. accountKey: storageKey,
  269. dbPath: latestPath('messages', `agent-workbench-${storageKey}.db`),
  270. agent: {
  271. claudeSessionFile: latestPath('messages', `claude-code-sessions-${storageKey}.json`),
  272. claudeProjectId: `${crypto.createHash('sha256').update(PROJECT_ROOT).digest('hex').slice(0, 12)}-${storageKey.slice(0, 8)}`,
  273. },
  274. qiwei: {
  275. uid: String(input.uid || '').trim(),
  276. guid: String(input.guid || '').trim(),
  277. userId: String(input.userId || '').trim(),
  278. nickname: String(input.nickname || '').trim(),
  279. corpName: String(input.corpName || '').trim(),
  280. apiBase: String(input.apiBase || '').trim(),
  281. ...(input.useConfiguredAllowlist === true ? {} : { allowedSenders: [], manualAllowedSenders: [] }),
  282. },
  283. };
  284. }
  285. function backfillCustomerIntelligence(db) {
  286. const version = '4';
  287. if (db.getSetting('customer_intelligence_backfill_version', '') === version) return { skipped: true, profileFields: 0, taskCount: 0, alertCount: 0 };
  288. const removed = db.lastIntelligenceMigration || { tasks: { removed: 0 }, alerts: { removed: 0 } };
  289. let profileFields = 0;
  290. let taskCount = 0;
  291. let alertCount = 0;
  292. for (const conversation of db.listConversations()) {
  293. const existing = db.getProfile(conversation.id);
  294. let current = { profile: { ...existing.profile }, tags: existing.tags || [] };
  295. for (const message of db.listMessages(conversation.id, 200).filter(item => item.direction === 'inbound')) {
  296. const intelligence = extractExplicitCustomerIntelligence(message.content, current.profile || {}, {});
  297. if (Object.keys(intelligence.profileUpdates).length) {
  298. current = { profile: { ...current.profile, ...intelligence.profileUpdates }, tags: current.tags };
  299. profileFields += Object.keys(intelligence.profileUpdates).length;
  300. }
  301. const changed = Object.keys(intelligence.profileUpdates).length > 0;
  302. const tasks = intelligence.tasks.map(item => ({ ...item, sourceMessageId: changed ? message.id : null }));
  303. const alerts = intelligence.alerts.map(item => ({ ...item, sourceMessageId: changed || item.managedBy === 'event' ? message.id : null }));
  304. taskCount += db.reconcileCustomerTasks(conversation.id, tasks, message.id).tasks.length;
  305. alertCount += db.reconcileCustomerAlerts(conversation.id, alerts, message.id).alerts.length;
  306. }
  307. if (Object.keys(current.profile).length) db.updateProfile(conversation.id, { ...existing.profile, ...current.profile }, existing.tags);
  308. }
  309. db.setSetting('customer_intelligence_backfill_version', version);
  310. if (profileFields || taskCount || alertCount) {
  311. db.audit({ actor: 'migration', action: 'customer_intelligence_backfilled', detail: { version, removed, profileFields, taskCount, alertCount } });
  312. }
  313. return { version, removed, profileFields, taskCount, alertCount };
  314. }
  315. function backfillCustomerMemory(service) {
  316. const version = '2';
  317. const db = service?.db;
  318. if (!db || !service.memory || db.getSetting('customer_memory_backfill_version', '') === version) {
  319. return { skipped: true, conversations: 0, captured: 0 };
  320. }
  321. let conversations = 0;
  322. let captured = 0;
  323. for (const conversation of db.listConversations()) {
  324. const result = service.memory.backfillConversation(conversation);
  325. conversations += 1;
  326. captured += Number(result.captured || 0);
  327. }
  328. db.setSetting('customer_memory_backfill_version', version);
  329. db.audit({ actor: 'migration', action: 'customer_memory_backfilled', detail: { version, conversations, captured } });
  330. return { version, conversations, captured };
  331. }
  332. const delay = ms => new Promise(resolve => setTimeout(resolve, ms));
  333. class QiweiAgentPoller {
  334. constructor({ config, db, qiwei, service }) {
  335. this.config = config;
  336. this.db = db;
  337. this.qiwei = qiwei;
  338. this.service = service;
  339. this.running = false;
  340. this.loopPromise = null;
  341. this.startedAt = 0;
  342. this.lastError = '';
  343. this.wakeLoop = null;
  344. }
  345. status() {
  346. return {
  347. running: this.running,
  348. syncKey: Number(this.db.getPollState('sync_key', '0')),
  349. lastError: this.lastError,
  350. startedAt: this.startedAt || null,
  351. };
  352. }
  353. async start() {
  354. if (this.running) return this.status();
  355. if (!this.qiwei.isConfigured()) throw new Error('Fmode 企微网关尚未配置,请先完成 Fmode 鉴权和企微扫码登录');
  356. if (this.config.allowedSenders.length === 0 && this.db.intakePolicy().mode === 'allowlist_only') throw new Error('企微联系人白名单为空,拒绝启动');
  357. await this.qiwei.checkLogin();
  358. this.running = true;
  359. this.startedAt = Math.floor(Date.now() / 1000);
  360. this.lastError = '';
  361. let syncKey = Number(this.db.getPollState('sync_key', '0')) || 0;
  362. if (syncKey === 0) syncKey = await this.establishBaseline();
  363. this.db.audit({ actor: 'runtime', action: 'poller_started', detail: { syncKey, allowlistCount: this.config.allowedSenders.length } });
  364. this.loopPromise = this.loop(syncKey);
  365. return this.status();
  366. }
  367. stop() {
  368. if (!this.running) return this.status();
  369. this.running = false;
  370. if (this.wakeLoop) this.wakeLoop();
  371. this.db.audit({ actor: 'runtime', action: 'poller_stopped', detail: this.status() });
  372. return this.status();
  373. }
  374. waitInterval() {
  375. return new Promise(resolve => {
  376. const timer = setTimeout(() => {
  377. this.wakeLoop = null;
  378. resolve();
  379. }, this.config.intervalMs);
  380. this.wakeLoop = () => {
  381. clearTimeout(timer);
  382. this.wakeLoop = null;
  383. resolve();
  384. };
  385. });
  386. }
  387. async establishBaseline() {
  388. let cursor = 0;
  389. let reachedEnd = false;
  390. let pages = 0;
  391. let count = 0;
  392. while (pages < this.config.initialSyncMaxPages) {
  393. const result = await this.qiwei.syncMessages(cursor, this.config.initialSyncLimit);
  394. const list = result.syncMsgList || [];
  395. pages += 1;
  396. count += list.length;
  397. const seqs = list.map(item => Number(item.seq)).filter(Number.isFinite);
  398. const next = Math.max(cursor, Number(result.travelSyncKey) || 0, seqs.length ? Math.max(...seqs) : 0);
  399. if (list.length === 0) { reachedEnd = true; break; }
  400. if (next <= cursor) throw new Error(`历史消息游标没有前进(seq=${cursor})`);
  401. cursor = next;
  402. }
  403. if (!reachedEnd) throw new Error(`历史消息超过 ${this.config.initialSyncMaxPages} 页,拒绝启动自动处理`);
  404. this.db.setPollState('sync_key', cursor);
  405. this.db.audit({ actor: 'runtime', action: 'poller_baseline_established', detail: { cursor, pages, skippedHistory: count } });
  406. return cursor;
  407. }
  408. async loop(initialSyncKey) {
  409. let syncKey = initialSyncKey;
  410. while (this.running) {
  411. try {
  412. const result = await this.qiwei.syncMessages(syncKey, 50);
  413. const list = result.syncMsgList || [];
  414. const seqs = list.map(item => Number(item.seq)).filter(Number.isFinite);
  415. const next = Math.max(syncKey, Number(result.travelSyncKey) || 0, seqs.length ? Math.max(...seqs) : 0);
  416. for (const message of list) await this.process(message);
  417. syncKey = next;
  418. this.db.setPollState('sync_key', syncKey);
  419. this.lastError = '';
  420. } catch (error) {
  421. this.lastError = error.message;
  422. this.db.audit({ actor: 'runtime', action: 'poller_error', detail: { message: error.message } });
  423. }
  424. if (this.running) await this.waitInterval();
  425. }
  426. }
  427. async process(message) {
  428. const allowlist = refreshAllowedSendersFromEnv(this.config);
  429. if (allowlist.changed) {
  430. this.db.audit({ actor: 'runtime', action: 'poller_allowlist_reloaded', detail: { count: allowlist.count } });
  431. }
  432. return ingestMessageForWorkbench({ config: this.config, db: this.db, service: this.service, qiwei: this.qiwei }, message, 'polling');
  433. }
  434. }
  435. const OUTBOUND_NOTIFICATION_MSG_TYPE = 2001;
  436. const OUTBOUND_CONTENT_MSG_TYPES = new Set([0, 1, 2]);
  437. function outboundContactId(message = {}, config = {}) {
  438. const senderId = String(message.senderId || '');
  439. const receiverId = String(message.receiverId || '');
  440. const allowedSenders = new Set((config.allowedSenders || []).map(String));
  441. const selfUserIds = new Set([
  442. String(config.selfUserId || ''),
  443. String(config.userId || ''),
  444. String(message.userId || ''),
  445. ].filter(Boolean));
  446. return selfUserIds.has(senderId) && allowedSenders.has(receiverId) ? receiverId : '';
  447. }
  448. function recordOutboundMessage(target, message, source, contactId) {
  449. const content = messageContent(message);
  450. if (!content) throw new Error('Outbound message content is empty');
  451. const timestamp = messageTimestamp(message.timestamp);
  452. if (!timestamp) throw new Error('Outbound message timestamp is invalid');
  453. const existing = target.db.getConversationByContactId(contactId);
  454. const conversation = target.db.ensureConversation(contactId, existing?.contact_name || message.receiverName || '白名单企微联系人');
  455. const inserted = target.db.insertMessage({
  456. conversationId: conversation.id,
  457. externalId: String(message.msgServerId || message.msgUniqueIdentifier || `${message.senderId}:${message.seq}`),
  458. direction: 'outbound',
  459. senderType: 'human',
  460. content,
  461. status: 'sent',
  462. createdAt: new Date(timestamp * 1000).toISOString(),
  463. raw: { ...message, source },
  464. });
  465. if (inserted.created) {
  466. target.db.audit({
  467. actor: 'runtime',
  468. action: 'outbound_message_captured',
  469. conversationId: conversation.id,
  470. entityId: inserted.message.id,
  471. detail: { source, seq: Number(message.seq) || 0 },
  472. });
  473. target.service?.emit?.('change', { type: 'message', conversationId: conversation.id });
  474. }
  475. return { status: inserted.created ? 'outbound_recorded' : 'outbound_duplicate', message: inserted.message };
  476. }
  477. async function syncOutboundNotification(target, notification, source, contactId) {
  478. if (!target.qiwei?.syncMessages) throw new Error('Outbound message sync is unavailable');
  479. const notificationSeq = Number(notification.seq) || 0;
  480. const senderId = String(notification.senderId || '');
  481. const delays = [0, 250, 750, 1500];
  482. for (const delayMs of delays) {
  483. if (delayMs) await new Promise(resolve => setTimeout(resolve, delayMs));
  484. const result = await target.qiwei.syncMessages(Math.max(0, notificationSeq - 1), 10);
  485. const message = (result.syncMsgList || []).find(item => (
  486. Number(item.seq) >= notificationSeq
  487. && String(item.senderId || '') === senderId
  488. && String(item.receiverId || '') === contactId
  489. && OUTBOUND_CONTENT_MSG_TYPES.has(Number(item.msgType))
  490. && Boolean(messageContent(item))
  491. ));
  492. if (message) return recordOutboundMessage(target, message, `${source}:outbound-sync`, contactId);
  493. }
  494. throw new Error(`Outbound message content is not available after notification seq=${notificationSeq}`);
  495. }
  496. async function ingestMessageForWorkbench(target, message = {}, source = 'callback') {
  497. if (isGroupMessage(message)) {
  498. return { status: 'ignored_group_reply_disabled' };
  499. }
  500. const contactId = outboundContactId(message, target.config);
  501. if (contactId && Number(message.msgType) === OUTBOUND_NOTIFICATION_MSG_TYPE) {
  502. return syncOutboundNotification(target, message, source, contactId);
  503. }
  504. if (contactId && OUTBOUND_CONTENT_MSG_TYPES.has(Number(message.msgType)) && messageContent(message)) {
  505. return recordOutboundMessage(target, message, source, contactId);
  506. }
  507. const intakePolicy = target.db.intakePolicy();
  508. const candidate = evaluatePolledMessage(message, { ...target.config, intakeMode: intakePolicy.mode });
  509. if (!candidate.eligible) return { status: `ignored_${candidate.reason || 'ineligible'}` };
  510. const { content, senderId, timestamp } = candidate;
  511. let onboarding = false;
  512. let allowAutoSend;
  513. if (candidate.unknownContact) {
  514. if (intakePolicy.mode === 'allowlist_only') return { status: 'ignored_not_allowlisted' };
  515. if (intakePolicy.mode === 'auto_enroll_autopilot' && intakePolicy.autopilotConfirmation !== INTAKE_AUTOPILOT_CONFIRMATION_VERSION) {
  516. target.db.audit({ actor: 'policy', action: 'auto_enroll_blocked_unconfirmed', detail: { contactId: senderId, source } });
  517. return { status: 'ignored_auto_enroll_unconfirmed' };
  518. }
  519. const autoIds = target.db.addAutoEnrolledContactId(senderId);
  520. target.config.autoEnrolledSenders = autoIds;
  521. target.config.allowedSenders = normalizeAllowlistIds([...(target.config.manualAllowedSenders || []), ...autoIds]);
  522. const conversation = target.db.ensureConversation(senderId, message.senderName || '企微客户');
  523. const effectiveMode = intakePolicy.mode === 'auto_enroll_review' ? 'review' : 'autopilot';
  524. target.service.setConversationMode(conversation.id, effectiveMode, 'policy:auto-enroll');
  525. const welcomeState = !intakePolicy.welcomeEnabled
  526. ? 'skipped'
  527. : (intakePolicy.mode === 'auto_enroll_autopilot' && intakePolicy.welcomeSendMode === 'send' ? 'send_pending' : 'draft_pending');
  528. const accountKey = String(target.config.accountKey || 'default');
  529. const existing = target.db.getOnboarding(accountKey, senderId);
  530. target.db.upsertOnboarding({
  531. accountKey,
  532. contactId: senderId,
  533. conversationId: conversation.id,
  534. contactName: message.senderName || '企微客户',
  535. source,
  536. policy: intakePolicy.mode,
  537. welcomeState,
  538. welcomeText: intakePolicy.welcomeText,
  539. });
  540. if (!existing) {
  541. target.db.audit({ actor: 'policy', action: 'contact_auto_enrolled', conversationId: conversation.id, detail: { contactId: senderId, source, policy: intakePolicy.mode, effectiveConversationMode: effectiveMode } });
  542. }
  543. onboarding = true;
  544. allowAutoSend = intakePolicy.mode !== 'auto_enroll_review';
  545. }
  546. return target.service.ingestInbound({
  547. externalId: String(message.msgServerId || message.msgUniqueIdentifier || `${senderId}:${message.seq}`),
  548. contactId: senderId,
  549. contactName: message.senderName || '企微客户',
  550. content,
  551. timestamp: new Date(timestamp * 1000).toISOString(),
  552. raw: {
  553. ...message,
  554. source,
  555. fromRoomId: message.fromRoomId || null,
  556. },
  557. }, { onboarding, allowAutoSend });
  558. }
  559. function createWorkbench(overrides = {}) {
  560. const config = loadAgentConfig(overrides.config || {});
  561. const db = overrides.db || new AgentWorkbenchDb(config.dbPath, {
  562. globalPaused: config.globalDefaultPaused,
  563. defaultMode: config.conversationDefaultMode,
  564. autoSendConfidence: config.autoSendConfidence,
  565. personalIntakeMode: config.intake.mode,
  566. welcomeEnabled: config.intake.welcomeEnabled,
  567. welcomeText: config.intake.welcomeText,
  568. welcomeSendMode: config.intake.welcomeSendMode,
  569. });
  570. hydrateAccountAllowlist(config, db);
  571. const staleWelcomeCount = db.markStaleWelcomeSending(
  572. config.accountKey,
  573. new Date(Date.now() - 5 * 60 * 1000).toISOString(),
  574. );
  575. if (staleWelcomeCount) {
  576. db.audit({ actor: 'runtime', action: 'onboarding_welcome_delivery_unknown', detail: { count: staleWelcomeCount } });
  577. }
  578. if (!overrides.db && !overrides.skipLegacyImport) {
  579. try { db.importCompatibleDatabase(config.legacyDbPath); } catch (error) {
  580. db.audit({ actor: 'migration', action: 'legacy_workbench_import_failed', detail: { message: error.message } });
  581. }
  582. const removed = db.cleanupInboundContentDuplicates(60);
  583. if (removed) db.audit({ actor: 'migration', action: 'duplicate_messages_cleaned', detail: { removed } });
  584. backfillCustomerIntelligence(db);
  585. }
  586. const knowledge = overrides.knowledge || new AgentKnowledgeStore({
  587. knowledgeDir: config.knowledgeDir,
  588. contextFiles: config.contextFiles,
  589. contextCharLimit: config.contextCharLimit,
  590. });
  591. backfillSentVoiceAudioPaths(db);
  592. const qiwei = overrides.qiwei || new FmodeQiweiClient(config.qiwei);
  593. const voice = overrides.voice || new VoiceCloneService({ config: config.voice, qiwei });
  594. const agent = overrides.agent || new QiweiAgentRuntime({ config: config.agent, knowledge });
  595. const service = overrides.service || new AgentWorkbenchService({ db, agent, qiwei, config });
  596. backfillCustomerMemory(service);
  597. const poller = overrides.poller || new QiweiAgentPoller({ config: config.qiwei, db, qiwei, service });
  598. return { config, db, knowledge, qiwei, voice, agent, service, poller };
  599. }
  600. function configuredStartupAccount() {
  601. return {
  602. uid: value('QIWEI_UID') || value('QIWE_UID'),
  603. guid: value('QIWEI_GUID') || value('QIWE_GUID'),
  604. apiBase: value('QIWEI_API_BASE') || value('QIWE_API_BASE'),
  605. userId: '',
  606. nickname: '',
  607. corpName: '',
  608. useConfiguredAllowlist: true,
  609. };
  610. }
  611. const startupAccount = configuredStartupAccount();
  612. let workbench = createWorkbench(startupAccount.uid ? {
  613. config: accountWorkbenchOverrides(startupAccount),
  614. skipLegacyImport: true,
  615. } : {});
  616. function migrateStartupWorkbench(target) {
  617. if (!startupAccount.uid || target.db.listConversations().length) return { imported: false, reason: 'target_not_empty' };
  618. const candidates = [latestPath('messages', 'agent-workbench.db'), target.config.legacyDbPath];
  619. for (const sourcePath of candidates) {
  620. try {
  621. const result = target.db.importCompatibleDatabase(sourcePath);
  622. if (!result.imported) continue;
  623. const removed = target.db.cleanupInboundContentDuplicates(60);
  624. backfillCustomerIntelligence(target.db);
  625. target.db.audit({
  626. actor: 'migration',
  627. action: 'account_workbench_migration_completed',
  628. detail: { source: path.basename(sourcePath), removedDuplicates: removed },
  629. });
  630. return result;
  631. } catch (error) {
  632. target.db.audit({
  633. actor: 'migration',
  634. action: 'account_workbench_migration_failed',
  635. detail: { source: path.basename(sourcePath), message: error.message },
  636. });
  637. }
  638. }
  639. return { imported: false, reason: 'source_missing_or_incompatible' };
  640. }
  641. migrateStartupWorkbench(workbench);
  642. let accountStatusCache = { checkedAt: 0, value: null };
  643. let accountStatusRefresh = null;
  644. let accountStatusRefreshKey = '';
  645. const accountLastOnlineAt = new Map();
  646. const accountOfflineChecks = new Map();
  647. const ONLINE_STATUS_GRACE_MS = 60000;
  648. const workbenches = new Map();
  649. function activeAccountMetadata() {
  650. const context = workbench.qiwei.context();
  651. return {
  652. uid: String(context.uid || workbench.config.qiwei.uid || '').trim(),
  653. guid: String(context.guid || workbench.config.qiwei.guid || '').trim(),
  654. apiBase: String(context.apiBase || workbench.config.qiwei.apiBase || '').trim(),
  655. userId: String(workbench.config.qiwei.userId || '').trim(),
  656. nickname: String(workbench.config.qiwei.nickname || '').trim(),
  657. corpName: String(workbench.config.qiwei.corpName || '').trim(),
  658. };
  659. }
  660. function applyActiveAccountContext(account) {
  661. setActiveQiweiContext(account);
  662. Object.assign(workbench.config.qiwei, account);
  663. if (workbench.qiwei && workbench.qiwei.config) Object.assign(workbench.qiwei.config, account);
  664. }
  665. function provisionalAccountStatus(selected = activeAccountMetadata(), statusText = '正在检测账号状态') {
  666. return {
  667. uid: selected.uid,
  668. guid: selected.guid,
  669. userId: selected.userId,
  670. configured: workbench.qiwei.isConfigured(),
  671. online: false,
  672. nickname: selected.nickname || selected.userId || '当前企微账号',
  673. corpName: selected.corpName || '',
  674. statusCode: null,
  675. statusText,
  676. };
  677. }
  678. const initialAccount = activeAccountMetadata();
  679. workbenches.set(accountRuntimeKey(initialAccount), workbench);
  680. applyActiveAccountContext(initialAccount);
  681. async function switchActiveAccount(input = {}) {
  682. const requestedGuid = String(input.guid || '').trim();
  683. const account = {
  684. uid: String(input.uid || '').trim(),
  685. guid: requestedGuid === 'server-managed' ? '' : requestedGuid,
  686. apiBase: String(input.apiBase || '').trim(),
  687. userId: String(input.userId || '').trim(),
  688. nickname: String(input.nickname || input.userId || '').trim(),
  689. corpName: String(input.corpName || '').trim(),
  690. };
  691. if (!account.uid) throw new Error('该账号缺少 Fmode 设备 uid,请重新扫码绑定后再切换');
  692. const current = activeAccountMetadata();
  693. const currentKey = accountRuntimeKey(current);
  694. const nextKey = accountRuntimeKey(account);
  695. const accountChanged = currentKey !== nextKey;
  696. if (accountChanged && workbench.poller.status().running) workbench.poller.stop();
  697. let nextWorkbench = workbenches.get(nextKey);
  698. if (!nextWorkbench) {
  699. nextWorkbench = createWorkbench({
  700. config: accountWorkbenchOverrides(account),
  701. skipLegacyImport: true,
  702. });
  703. workbenches.set(nextKey, nextWorkbench);
  704. }
  705. workbench = nextWorkbench;
  706. const currentAccount = activeAccountMetadata();
  707. const accountPatch = {
  708. ...account,
  709. apiBase: account.apiBase || currentAccount.apiBase,
  710. guid: account.guid || (accountChanged ? '' : currentAccount.guid),
  711. };
  712. applyActiveAccountContext({ ...currentAccount, ...accountPatch });
  713. saveQiweiClientConfig({
  714. uid: accountPatch.uid,
  715. guid: accountPatch.guid,
  716. apiBase: accountPatch.apiBase,
  717. envRoot: PROJECT_ROOT,
  718. });
  719. const status = provisionalAccountStatus();
  720. accountStatusCache = { checkedAt: Date.now(), value: status };
  721. void refreshAccountStatus();
  722. return {
  723. status: 'ok',
  724. assistantMessage: `已切换到账号:${status.nickname || account.nickname || account.userId || account.uid}`,
  725. summary: {
  726. switched: accountChanged,
  727. storageKey: nextKey,
  728. online: status.online,
  729. listenerStopped: accountChanged,
  730. },
  731. data: { account: status },
  732. };
  733. }
  734. async function refreshAccountStatus() {
  735. const selected = activeAccountMetadata();
  736. const targetWorkbench = workbench;
  737. const refreshKey = accountRuntimeKey(selected);
  738. if (accountStatusRefresh && accountStatusRefreshKey === refreshKey) return accountStatusRefresh;
  739. accountStatusRefreshKey = refreshKey;
  740. accountStatusRefresh = (async () => {
  741. let next;
  742. try {
  743. const data = await targetWorkbench.qiwei.checkLogin();
  744. const online = Number(data.userOnlineStatus) === 2 && Number(data.errorCode || 0) === 0;
  745. next = {
  746. uid: selected.uid,
  747. guid: selected.guid,
  748. userId: data.userId || selected.userId,
  749. configured: data.configured !== false,
  750. online,
  751. nickname: data.nickname || selected.nickname || selected.userId || '当前企微账号',
  752. corpName: data.corpName || selected.corpName || '',
  753. statusCode: data.userOnlineStatus ?? null,
  754. statusText: online ? (data.statusFallback ? '账号在线(状态接口待复核)' : '账号在线') : '账号离线',
  755. };
  756. if (online) {
  757. accountLastOnlineAt.set(refreshKey, Date.now());
  758. accountOfflineChecks.set(refreshKey, 0);
  759. } else {
  760. const offlineChecks = Number(accountOfflineChecks.get(refreshKey) || 0) + 1;
  761. accountOfflineChecks.set(refreshKey, offlineChecks);
  762. const lastOnlineAt = Number(accountLastOnlineAt.get(refreshKey) || 0);
  763. if (offlineChecks < 2 && lastOnlineAt && Date.now() - lastOnlineAt < ONLINE_STATUS_GRACE_MS) {
  764. next.online = true;
  765. next.statusCode = 2;
  766. next.statusText = '账号在线(正在复核)';
  767. }
  768. }
  769. } catch {
  770. next = {
  771. ...provisionalAccountStatus(selected, '状态检测失败'),
  772. configured: targetWorkbench.qiwei.isConfigured(),
  773. };
  774. const lastOnlineAt = Number(accountLastOnlineAt.get(refreshKey) || 0);
  775. if (lastOnlineAt && Date.now() - lastOnlineAt < ONLINE_STATUS_GRACE_MS) {
  776. next.online = true;
  777. next.statusCode = 2;
  778. next.statusText = '账号在线(状态刷新中)';
  779. }
  780. }
  781. if (accountRuntimeKey(activeAccountMetadata()) === refreshKey) {
  782. accountStatusCache = { checkedAt: Date.now(), value: next };
  783. }
  784. if (next.userId) {
  785. targetWorkbench.config.qiwei.userId = String(next.userId);
  786. targetWorkbench.config.qiwei.selfUserId = String(next.userId);
  787. }
  788. return next;
  789. })().finally(() => {
  790. if (accountStatusRefreshKey === refreshKey) {
  791. accountStatusRefresh = null;
  792. accountStatusRefreshKey = '';
  793. }
  794. });
  795. return accountStatusRefresh;
  796. }
  797. async function detectAccountStatus(force = false) {
  798. const selected = activeAccountMetadata();
  799. const selectedKey = accountRuntimeKey(selected);
  800. const cacheMatches = accountStatusCache.value && accountRuntimeKey(accountStatusCache.value) === selectedKey;
  801. if (force) return refreshAccountStatus();
  802. if (cacheMatches) {
  803. if (Date.now() - accountStatusCache.checkedAt >= 8000) void refreshAccountStatus();
  804. return accountStatusCache.value;
  805. }
  806. const provisional = provisionalAccountStatus(selected);
  807. accountStatusCache = { checkedAt: Date.now(), value: provisional };
  808. void refreshAccountStatus();
  809. return provisional;
  810. }
  811. function maskedId(value) {
  812. const text = String(value || '');
  813. if (text.length <= 4) return '测试联系人';
  814. return `${text.slice(0, 2)}***${text.slice(-2)}`;
  815. }
  816. function completeness(profile = {}) {
  817. const values = [
  818. profile.preferredRegion || profile.region || profile.district || profile.intent_area || profile.districts,
  819. profile.budgetWan || profile.budget || profile.budgetMax || profile.budget_max,
  820. profile.layout || profile.rooms || profile.house_type,
  821. profile.area || profile.areaMin || profile.area_min,
  822. profile.decoration,
  823. profile.timeline || profile.urgency,
  824. ];
  825. return Math.round(values.filter(value => Array.isArray(value) ? value.length : Boolean(value)).length / values.length * 100);
  826. }
  827. function parseJson(value, fallback) {
  828. try { return value ? JSON.parse(value) : fallback; }
  829. catch { return fallback; }
  830. }
  831. function conversationChannelInfo(row, messages = []) {
  832. const contactId = String(row?.contact_id || '').trim();
  833. let roomId = '';
  834. for (const message of messages) {
  835. const raw = parseJson(message.raw_json, {});
  836. if (!isGroupMessage(raw)) continue;
  837. roomId = roomIdOf(raw) || contactId;
  838. break;
  839. }
  840. if (!roomId && /(?:@chatroom$|^(?:room|group|chatroom|r[-_:]))/i.test(contactId)) roomId = contactId;
  841. const channelType = roomId ? 'group' : 'private';
  842. return {
  843. channelType,
  844. channelLabel: channelType === 'group' ? '客户群聊' : '客户私聊',
  845. channelId: roomId || contactId,
  846. };
  847. }
  848. function publicConversation(row) {
  849. const detail = workbench.service.conversationDetail(row.id);
  850. const claudeSession = getCustomerSessionGuide(row, {
  851. sessionFile: workbench.config.agent.claudeSessionFile,
  852. });
  853. const drafts = detail.drafts || [];
  854. const latestInbound = [...(detail.messages || [])].reverse().find(message => message.direction === 'inbound') || null;
  855. const currentDrafts = latestInbound
  856. ? drafts.filter(item => item.inbound_message_id === latestInbound.id)
  857. : [];
  858. const pending = currentDrafts.find(item => item.status === 'pending') || null;
  859. const latestDraft = pending || currentDrafts.find(item => ['sent', 'approved'].includes(item.status)) || null;
  860. const currentAgentOutcome = latestInbound && detail.agentOutcome?.entityId === latestInbound.id ? detail.agentOutcome : null;
  861. const agentError = ['agent_failed', 'agent_not_configured', 'autopilot_send_failed'].includes(currentAgentOutcome?.action)
  862. ? { ...currentAgentOutcome, ...friendlyAgentError(currentAgentOutcome.message) }
  863. : null;
  864. const agentNotice = currentAgentOutcome?.action === 'agent_no_reply_needed' ? currentAgentOutcome : null;
  865. const rawProfile = detail.profile?.profile || {};
  866. const { __evidence: profileEvidence = {}, ...profile } = rawProfile;
  867. const customerTasks = detail.tasks || [];
  868. const customerAlerts = detail.alerts || [];
  869. const citations = latestDraft?.citations || [];
  870. const cutoverAt = Date.parse(workbench.db.getSetting('agent_cutover_at', '')) || Date.now();
  871. const displayMessages = [...new Map((detail.messages || []).map(message => [message.id, message])).values()]
  872. .sort((a, b) => Date.parse(a.created_at) - Date.parse(b.created_at))
  873. .slice(-60);
  874. const channel = conversationChannelInfo(row, displayMessages);
  875. const visibleEntityIds = new Set(displayMessages.map(message => message.id));
  876. const visibleAudit = (detail.audit || []).filter(item =>
  877. Date.parse(item.created_at) >= cutoverAt || visibleEntityIds.has(item.entity_id)
  878. ).slice(0, 40).map(item => {
  879. const publicDetail = { ...(item.detail || {}) };
  880. delete publicDetail.rawMessage;
  881. return { ...item, detail_json: undefined, detail: publicDetail };
  882. });
  883. return {
  884. id: row.id,
  885. displayName: row.contact_name || '白名单测试联系人',
  886. maskedId: maskedId(row.contact_id),
  887. ...channel,
  888. mode: row.mode,
  889. source: 'live',
  890. claudeSession,
  891. messages: displayMessages.map(message => {
  892. const raw = parseJson(message.raw_json, {});
  893. if (Number(raw.msgType) === 16 && raw.audioPath) {
  894. raw.audioAvailable = true;
  895. delete raw.audioPath;
  896. }
  897. return {
  898. id: message.id,
  899. role: message.direction === 'inbound' ? 'customer' : message.sender_type,
  900. content: message.content,
  901. timestamp: message.created_at,
  902. status: message.status,
  903. source: message.direction === 'inbound' ? 'live' : message.sender_type,
  904. raw,
  905. };
  906. }),
  907. analysis: {
  908. intent: latestDraft?.intent || '',
  909. intentLabel: latestDraft?.intent || (agentError ? 'Agent 上游不可用' : agentNotice ? '无需回复' : '待 Agent 处理'),
  910. demand: profile,
  911. completenessScore: completeness(profile),
  912. knowledgeSources: citations.length
  913. ? citations.map(item => `${item.heading || item.source}${item.source ? ` · ${item.source}` : ''}`)
  914. : ['真实企微消息', '客户画像', '企业规则库与知识库'],
  915. reasoning: latestDraft?.reason || agentError?.message || agentNotice?.message || '消息已进入真实企微链路,等待 Agent 生成可审核草稿。',
  916. },
  917. customerIntelligence: {
  918. profile,
  919. profileEvidence,
  920. profileUpdatedAt: detail.profile?.updatedAt || null,
  921. tags: detail.profile?.tags || [],
  922. tasks: customerTasks.map(item => ({
  923. id: item.id,
  924. businessKey: item.business_key,
  925. type: item.type,
  926. title: item.title,
  927. owner: item.owner,
  928. dueAt: item.due_at,
  929. priority: item.priority,
  930. status: item.status,
  931. reason: item.reason,
  932. evidence: item.evidence,
  933. evidenceItems: parseJson(item.evidence_json, []),
  934. resolutionReason: item.resolution_reason,
  935. officialTodoId: item.official_todo_id,
  936. officialSyncStatus: item.official_sync_status,
  937. officialSyncedAt: item.official_synced_at,
  938. updatedAt: item.updated_at,
  939. })),
  940. alerts: customerAlerts.map(item => ({
  941. id: item.id,
  942. businessKey: item.business_key,
  943. type: item.type,
  944. severity: item.severity,
  945. title: item.title,
  946. detail: item.detail,
  947. evidence: item.evidence,
  948. evidenceItems: parseJson(item.evidence_json, []),
  949. resolutionReason: item.resolution_reason,
  950. recommendedAction: item.recommended_action,
  951. status: item.status,
  952. updatedAt: item.updated_at,
  953. })),
  954. memories: (detail.memories || []).map(item => ({
  955. id: item.id,
  956. type: item.type,
  957. content: item.content,
  958. confidence: item.confidence,
  959. importance: item.importance,
  960. status: item.status,
  961. sourceMessageIds: item.source_message_ids || [],
  962. direction: item.direction,
  963. createdBy: item.created_by,
  964. expiresAt: item.expires_at,
  965. updatedAt: item.updated_at,
  966. })),
  967. memorySnapshot: detail.memorySnapshot ? {
  968. version: detail.memorySnapshot.version,
  969. coreChars: String(detail.memorySnapshot.compact_text || '').length,
  970. createdAt: detail.memorySnapshot.created_at,
  971. } : null,
  972. summary: {
  973. openTasks: customerTasks.filter(item => ['open', 'in_progress'].includes(item.status)).length,
  974. openAlerts: customerAlerts.filter(item => item.status === 'open').length,
  975. highAlerts: customerAlerts.filter(item => item.status === 'open' && ['high', 'critical'].includes(item.severity)).length,
  976. activeMemories: (detail.memories || []).filter(item => item.status === 'active').length,
  977. },
  978. },
  979. pendingReply: pending ? {
  980. id: pending.id,
  981. content: pending.content,
  982. status: pending.status,
  983. confidence: pending.confidence,
  984. reason: pending.reason,
  985. requiresHuman: pending.requires_human,
  986. citations: pending.citations,
  987. toolTrace: pending.tool_trace,
  988. createdAt: pending.created_at,
  989. } : null,
  990. drafts,
  991. audit: visibleAudit,
  992. agentError,
  993. agentNotice,
  994. onboarding: detail.onboarding ? {
  995. state: detail.onboarding.welcome_state,
  996. attemptCount: detail.onboarding.attempt_count,
  997. lastError: detail.onboarding.last_error || '',
  998. lastAttemptAt: detail.onboarding.last_attempt_at || null,
  999. sentAt: detail.onboarding.sent_at || null,
  1000. } : null,
  1001. lastMessageAt: maxConversationTimestamp(displayMessages.at(-1)?.created_at, row.last_message_at),
  1002. updatedAt: row.updated_at,
  1003. };
  1004. }
  1005. async function getAgentStatus() {
  1006. const account = await detectAccountStatus();
  1007. const state = workbench.service.state(workbench.poller.status());
  1008. const product = getProductMode();
  1009. const listener = listenerStateWithRuntime(product, state.poller);
  1010. const daemonCallback = getRelayDaemonStatus();
  1011. const runtimeRelayRunning = product.mode === 'enterprise'
  1012. && listener.transport === 'server_relay'
  1013. && listener.running;
  1014. const callback = runtimeRelayRunning
  1015. ? {
  1016. ...daemonCallback,
  1017. running: true,
  1018. processAlive: true,
  1019. pid: listener.runtimePid || daemonCallback.pid,
  1020. heartbeatAt: readRuntimeState().updatedAt || daemonCallback.heartbeatAt,
  1021. }
  1022. : daemonCallback;
  1023. const ingress = product.mode === 'enterprise'
  1024. ? { mode: 'server_callback', running: callback.running, durableQueue: true, ...callback }
  1025. : { mode: 'local_polling', running: state.poller.running, durableQueue: false };
  1026. return {
  1027. status: 'ok',
  1028. data: {
  1029. globalMode: state.global.paused ? 'paused' : state.global.defaultMode,
  1030. global: state.global,
  1031. listener,
  1032. ingress,
  1033. callback,
  1034. account,
  1035. product,
  1036. agent: state.agent,
  1037. knowledge: workbench.knowledge.stats(),
  1038. voice: workbench.voice.status(),
  1039. config: {
  1040. allowedSenderCount: state.qiwei.allowlistCount,
  1041. testMode: false,
  1042. demoMode: false,
  1043. transport: state.qiwei.transport,
  1044. pollIntervalMs: workbench.config.qiwei.intervalMs,
  1045. personalIngress: intakePolicyPayload(workbench),
  1046. },
  1047. safety: {
  1048. whitelistEnabled: state.qiwei.allowlistCount > 0,
  1049. defaultReviewMode: true,
  1050. globalPauseSupported: true,
  1051. messageSendRequiresWhitelist: true,
  1052. demoReplyDisabled: true,
  1053. },
  1054. },
  1055. };
  1056. }
  1057. async function ingestWebhookMessage(message = {}) {
  1058. return ingestMessageForWorkbench({
  1059. config: workbench.config.qiwei,
  1060. db: workbench.db,
  1061. service: workbench.service,
  1062. qiwei: workbench.qiwei,
  1063. }, message, 'relay_callback');
  1064. }
  1065. async function getAllowlistCandidates() {
  1066. const selectedIds = [...workbench.config.qiwei.allowedSenders];
  1067. const candidates = new Map();
  1068. let remoteCount = 0;
  1069. let warning = '';
  1070. try {
  1071. const account = await detectAccountStatus(true);
  1072. if (account.online) {
  1073. const result = await workbench.qiwei.listExternalContacts({ limit: 500, maxPages: 10 });
  1074. remoteCount = Number(result.contactCount) || result.contacts.length;
  1075. for (const item of result.contacts) {
  1076. const contact = normalizeAllowlistContact(item);
  1077. if (contact) candidates.set(contact.id, contact);
  1078. }
  1079. } else {
  1080. warning = '当前企微账号离线,只显示已保存或已有会话中的联系人';
  1081. }
  1082. } catch {
  1083. warning = '联系人列表暂时读取失败,仍可管理已保存联系人或手工添加联系人 ID';
  1084. }
  1085. for (const row of workbench.db.listConversations()) {
  1086. const contact = normalizeAllowlistContact({ contactId: row.contact_id }, row.contact_name);
  1087. if (contact && !candidates.has(contact.id)) candidates.set(contact.id, contact);
  1088. }
  1089. for (const id of selectedIds) {
  1090. if (!candidates.has(id)) candidates.set(id, normalizeAllowlistContact({ id }, '已保存联系人'));
  1091. }
  1092. const contacts = [...candidates.values()]
  1093. .map(contact => ({ ...contact, selected: selectedIds.includes(contact.id) }))
  1094. .sort((a, b) => Number(b.selected) - Number(a.selected) || a.displayName.localeCompare(b.displayName, 'zh-CN'));
  1095. return {
  1096. status: 'ok',
  1097. assistantMessage: warning || `已读取 ${contacts.length} 位可选联系人`,
  1098. summary: { selectedCount: selectedIds.length, candidateCount: contacts.length, remoteCount },
  1099. data: { selectedIds, contacts },
  1100. warnings: warning ? [warning] : [],
  1101. errors: [],
  1102. };
  1103. }
  1104. function updateAllowlistForWorkbench(target, input = {}) {
  1105. refreshAllowedSendersFromEnv(target.config.qiwei);
  1106. const previousIds = new Set(target.config.qiwei.allowedSenders.map(String));
  1107. const ids = normalizeAllowlistIds(input.contactIds || input.allowedSenders || []);
  1108. const retainedAutomatic = target.db.getAutoEnrolledContactIds().filter(id => ids.includes(id));
  1109. target.db.setSetting('manual_allowed_sender_ids', JSON.stringify(ids));
  1110. target.db.setSetting('auto_enrolled_contact_ids', JSON.stringify(retainedAutomatic));
  1111. target.config.qiwei.manualAllowedSenders = [...ids];
  1112. target.config.qiwei.autoEnrolledSenders = retainedAutomatic;
  1113. target.config.qiwei.allowedSenders = normalizeAllowlistIds([...ids, ...retainedAutomatic]);
  1114. let listenerStopped = false;
  1115. const contactNames = input.contactNames || {};
  1116. const addedIds = ids.filter(id => !previousIds.has(id));
  1117. for (const contactId of addedIds) {
  1118. target.db.ensureConversation(contactId, String(contactNames[contactId] || ''));
  1119. }
  1120. if (ids.length && input.autoStart !== false) {
  1121. target.db.setSetting('listener_enabled', 'true');
  1122. target.service.setGlobal({ paused: false }, input.actor || 'allowlist');
  1123. } else if (!ids.length) {
  1124. target.db.setSetting('listener_enabled', 'false');
  1125. }
  1126. if (!ids.length && target.poller.status().running) {
  1127. target.poller.stop();
  1128. listenerStopped = true;
  1129. }
  1130. target.db.audit({
  1131. actor: input.actor || 'human',
  1132. action: 'allowlist_updated',
  1133. detail: { count: ids.length, addedCount: addedIds.length, listenerStopped, autoStart: ids.length > 0 && input.autoStart !== false },
  1134. });
  1135. return {
  1136. status: 'ok',
  1137. assistantMessage: ids.length
  1138. ? `白名单已保存,共 ${ids.length} 位联系人;AI 监听将自动保持开启并按当前审核策略处理`
  1139. : '白名单已清空,AI 监听已停止',
  1140. summary: { selectedCount: ids.length, addedCount: addedIds.length, listenerStopped, autoStart: ids.length > 0 && input.autoStart !== false },
  1141. data: { selectedIds: ids, selectedCount: ids.length, addedIds },
  1142. warnings: [],
  1143. errors: [],
  1144. };
  1145. }
  1146. function addAllowlistContacts(input = {}) {
  1147. refreshAllowedSendersFromEnv(workbench.config.qiwei);
  1148. const contactIds = normalizeAllowlistIds(input.contactIds || []);
  1149. return updateAllowlist({
  1150. ...input,
  1151. contactIds: [...workbench.config.qiwei.allowedSenders, ...contactIds],
  1152. autoStart: input.autoStart !== false,
  1153. actor: input.actor || 'agent:test-contact',
  1154. });
  1155. }
  1156. function intakePolicyPayload(target = workbench) {
  1157. const policy = target.db.intakePolicy();
  1158. const onboardings = target.db.listOnboardings(target.config.accountKey, { limit: 500 });
  1159. return {
  1160. mode: policy.mode,
  1161. welcomeEnabled: policy.welcomeEnabled,
  1162. welcomeText: policy.welcomeText,
  1163. welcomeSendMode: policy.welcomeSendMode,
  1164. effectiveConversationMode: policy.mode === 'auto_enroll_autopilot' ? 'autopilot' : 'review',
  1165. autopilotConfirmed: policy.autopilotConfirmation === INTAKE_AUTOPILOT_CONFIRMATION_VERSION,
  1166. onboarding: {
  1167. total: onboardings.length,
  1168. failed: onboardings.filter(item => item.welcome_state === 'failed').length,
  1169. deliveryUnknown: onboardings.filter(item => item.welcome_state === 'delivery_unknown').length,
  1170. },
  1171. };
  1172. }
  1173. function updateAllowlist(input = {}) {
  1174. return updateAllowlistForWorkbench(workbench, input);
  1175. }
  1176. function getIntakePolicy() {
  1177. return { status: 'ok', data: intakePolicyPayload(workbench), warnings: [], errors: [] };
  1178. }
  1179. function updateIntakePolicyForWorkbench(target, input = {}, actor = 'human') {
  1180. const current = target.db.intakePolicy();
  1181. const mode = input.mode === undefined ? current.mode : String(input.mode);
  1182. const welcomeSendMode = input.welcomeSendMode === undefined ? current.welcomeSendMode : String(input.welcomeSendMode);
  1183. const welcomeEnabled = input.welcomeEnabled === undefined ? current.welcomeEnabled : input.welcomeEnabled === true;
  1184. const welcomeText = input.welcomeText === undefined ? current.welcomeText : String(input.welcomeText || '').trim();
  1185. if (!['allowlist_only', 'auto_enroll_review', 'auto_enroll_autopilot'].includes(mode)) throw new Error('不支持的个人消息接入模式');
  1186. if (!['draft', 'send'].includes(welcomeSendMode)) throw new Error('不支持的欢迎语发送模式');
  1187. if (welcomeEnabled && !welcomeText) throw new Error('启用欢迎语时必须填写欢迎内容');
  1188. if (welcomeSendMode === 'send' && mode !== 'auto_enroll_autopilot') throw new Error('欢迎语直接发送仅可用于全自动新联系人接入');
  1189. if (mode === 'auto_enroll_autopilot' && input.confirmation !== INTAKE_AUTOPILOT_CONFIRMATION) {
  1190. throw new Error('开启全自动新联系人接入需要二次确认');
  1191. }
  1192. const confirmed = mode === 'auto_enroll_autopilot';
  1193. target.db.setIntakePolicy({
  1194. mode,
  1195. welcomeEnabled,
  1196. welcomeText,
  1197. welcomeSendMode,
  1198. autopilotConfirmation: confirmed ? INTAKE_AUTOPILOT_CONFIRMATION_VERSION : '',
  1199. autopilotConfirmedAt: confirmed ? new Date().toISOString() : '',
  1200. autopilotConfirmedBy: confirmed ? String(actor || 'human') : '',
  1201. });
  1202. target.db.audit({
  1203. actor,
  1204. action: 'personal_intake_policy_updated',
  1205. detail: { mode, welcomeEnabled, welcomeSendMode, effectiveConversationMode: mode === 'auto_enroll_autopilot' ? 'autopilot' : 'review' },
  1206. });
  1207. return intakePolicyPayload(target);
  1208. }
  1209. function updateIntakePolicy(input = {}) {
  1210. return {
  1211. status: 'ok',
  1212. assistantMessage: '个人消息接入策略已保存',
  1213. data: updateIntakePolicyForWorkbench(workbench, input, 'human'),
  1214. warnings: [],
  1215. errors: [],
  1216. };
  1217. }
  1218. async function retryOnboardingWelcome(identifier) {
  1219. const value = String(identifier || '');
  1220. const conversation = workbench.db.getConversation(value);
  1221. const contactId = conversation?.contact_id || value;
  1222. const result = await workbench.service.retryOnboardingWelcome(contactId, 'human');
  1223. return {
  1224. status: 'ok',
  1225. assistantMessage: result.status === 'welcome_sent' ? '欢迎语重试已发送' : '欢迎语重试已处理',
  1226. data: result,
  1227. };
  1228. }
  1229. function displayableConversations() {
  1230. return workbench.db.listConversations();
  1231. }
  1232. function getConversations() {
  1233. const conversations = sortConversationsByRecency(displayableConversations().map(publicConversation));
  1234. return {
  1235. status: 'ok',
  1236. data: {
  1237. conversations,
  1238. summary: {
  1239. privateCount: conversations.filter(item => item.channelType === 'private').length,
  1240. groupCount: conversations.filter(item => item.channelType === 'group').length,
  1241. },
  1242. },
  1243. };
  1244. }
  1245. function getResponseMonitor() {
  1246. const allowlist = new Set(workbench.config.qiwei.allowedSenders.map(String));
  1247. const conversations = displayableConversations()
  1248. .filter(row => allowlist.has(String(row.contact_id || '')))
  1249. .flatMap(row => {
  1250. const conversation = publicConversation(row);
  1251. if (conversation.channelType === 'group') return [];
  1252. return [{
  1253. id: conversation.id,
  1254. contactId: row.contact_id,
  1255. displayName: conversation.displayName,
  1256. messages: conversation.messages,
  1257. handledMessageId: conversation.agentNotice?.entityId || '',
  1258. }];
  1259. });
  1260. const account = activeAccountMetadata();
  1261. const data = responseMonitor({
  1262. projectRoot: PROJECT_ROOT,
  1263. conversations,
  1264. selfUserIds: [workbench.config.qiwei.selfUserId, account.userId],
  1265. selfNames: [account.nickname],
  1266. warningMinutes: workbench.config.qiwei.responseReminderMinutes,
  1267. urgentMinutes: workbench.config.qiwei.responseUrgentMinutes,
  1268. });
  1269. return {
  1270. status: 'ok',
  1271. assistantMessage: data.summary.overdue
  1272. ? `当前有 ${data.summary.overdue} 个会话超过回复时效,请客服优先处理。`
  1273. : '当前监听范围内没有超时未回复会话。',
  1274. summary: data.summary,
  1275. data,
  1276. warnings: [],
  1277. errors: [],
  1278. };
  1279. }
  1280. function updateCustomerProfile(conversationId, input = {}) {
  1281. const conversation = workbench.db.getConversation(conversationId);
  1282. if (!conversation) throw new Error('客户会话不存在');
  1283. const current = workbench.db.getProfile(conversationId);
  1284. const nextProfile = { ...(current.profile || {}) };
  1285. const evidence = { ...(nextProfile.__evidence || {}) };
  1286. const patch = input.profile && typeof input.profile === 'object' ? input.profile : {};
  1287. const changedFields = [];
  1288. for (const [field, rawValue] of Object.entries(patch)) {
  1289. if (!field || field.startsWith('__')) continue;
  1290. const value = typeof rawValue === 'string' ? rawValue.trim() : rawValue;
  1291. if (value === '' || value === null || value === undefined) delete nextProfile[field];
  1292. else nextProfile[field] = value;
  1293. evidence[field] = {
  1294. text: String(input.reason || '客户管理人工核对').trim(),
  1295. sourceMessageId: null,
  1296. source: 'human',
  1297. updatedAt: new Date().toISOString(),
  1298. };
  1299. changedFields.push(field);
  1300. }
  1301. nextProfile.__evidence = evidence;
  1302. const tags = input.tags === undefined
  1303. ? current.tags
  1304. : [...new Set((Array.isArray(input.tags) ? input.tags : String(input.tags || '').split(/[,,]/)).map(item => String(item).trim()).filter(Boolean))];
  1305. const updated = workbench.db.updateProfile(conversationId, nextProfile, tags);
  1306. const memory = workbench.service.memory.syncHumanProfile(conversationId, patch);
  1307. workbench.db.audit({
  1308. actor: 'human',
  1309. action: 'customer_profile_updated',
  1310. conversationId,
  1311. entityId: conversationId,
  1312. detail: { fields: changedFields, tagCount: tags.length, reason: String(input.reason || '').trim() },
  1313. });
  1314. const { __evidence, ...visibleProfile } = updated.profile || {};
  1315. return {
  1316. status: 'ok',
  1317. assistantMessage: `客户“${conversation.contact_name || '未命名客户'}”的主档已更新。`,
  1318. summary: { changedFields, tagCount: tags.length },
  1319. data: { profile: visibleProfile, profileEvidence: __evidence || {}, tags, updatedAt: updated.updatedAt, memory },
  1320. warnings: [],
  1321. errors: [],
  1322. };
  1323. }
  1324. function resolveConversationSyncScope(input = {}, allowlist = new Set(), db = workbench.db) {
  1325. const conversationId = String(input.conversationId || '').trim();
  1326. let contactId = String(input.contactId || '').trim();
  1327. if (conversationId) {
  1328. const conversation = db.getConversation(conversationId);
  1329. if (!conversation) throw new Error('当前会话不存在,请刷新后重试');
  1330. if (conversationChannelInfo(conversation, db.listMessages(conversation.id, 100)).channelType !== 'private') {
  1331. throw new Error('当前采集仅支持个人聊天会话');
  1332. }
  1333. contactId = String(conversation.contact_id || '').trim();
  1334. }
  1335. if (contactId && !allowlist.has(contactId)) throw new Error('当前联系人不在个人消息白名单中');
  1336. return {
  1337. conversationId,
  1338. contactId,
  1339. contacts: contactId ? new Set([contactId]) : allowlist,
  1340. scope: contactId ? 'conversation' : 'allowlist',
  1341. };
  1342. }
  1343. async function syncConversations(input = {}) {
  1344. const allowlist = new Set(workbench.config.qiwei.allowedSenders.map(String));
  1345. if (!allowlist.size) throw new Error('测试联系人白名单为空,无法同步会话');
  1346. const syncScope = resolveConversationSyncScope(input, allowlist);
  1347. const account = await detectAccountStatus(true);
  1348. if (!account.online) throw new Error('测试账号当前不在线,无法同步企微会话');
  1349. const grouped = new Map();
  1350. const seen = new Set();
  1351. const seenSemantic = new Set();
  1352. let cursor = 0;
  1353. let pages = 0;
  1354. let scannedMessages = 0;
  1355. while (pages < workbench.config.qiwei.initialSyncMaxPages) {
  1356. const result = await workbench.qiwei.syncMessages(cursor, workbench.config.qiwei.initialSyncLimit);
  1357. const list = Array.isArray(result.syncMsgList) ? result.syncMsgList : [];
  1358. scannedMessages += list.length;
  1359. for (const message of list) {
  1360. if (isGroupMessage(message)) continue;
  1361. const senderId = String(message.senderId || '');
  1362. const receiverId = String(message.receiverId || '');
  1363. const contactId = syncScope.contacts.has(senderId) ? senderId : syncScope.contacts.has(receiverId) ? receiverId : '';
  1364. const content = String(message.msgData?.content || '').trim();
  1365. if (!contactId || !content || ![0, 1, 2].includes(Number(message.msgType))) continue;
  1366. const inbound = senderId === contactId;
  1367. const timestampSeconds = messageTimestamp(message.timestamp);
  1368. if (!timestampSeconds) continue;
  1369. const timestamp = new Date(timestampSeconds * 1000).toISOString();
  1370. const externalId = String(message.msgServerId || message.msgUniqueIdentifier || '');
  1371. const dedupeKey = externalId || `${contactId}|${inbound ? 'in' : 'out'}|${timestamp}|${content}`;
  1372. if (seen.has(dedupeKey)) continue;
  1373. seen.add(dedupeKey);
  1374. const semanticKey = `${contactId}|${inbound ? 'in' : 'out'}|${timestamp}|${content}`;
  1375. if (seenSemantic.has(semanticKey)) continue;
  1376. seenSemantic.add(semanticKey);
  1377. if (!grouped.has(contactId)) grouped.set(contactId, { contactName: '', messages: [] });
  1378. const group = grouped.get(contactId);
  1379. if (inbound && message.senderName) group.contactName = String(message.senderName);
  1380. group.messages.push({
  1381. externalId: externalId || null,
  1382. inbound,
  1383. content,
  1384. timestamp,
  1385. raw: { seq: message.seq, msgType: message.msgType, timestamp: message.timestamp, source: 'manual_sync' },
  1386. });
  1387. }
  1388. const seqs = list.map(item => Number(item.seq)).filter(Number.isFinite);
  1389. const next = Math.max(cursor, Number(result.travelSyncKey) || 0, seqs.length ? Math.max(...seqs) : 0);
  1390. pages += 1;
  1391. if (!list.length || next <= cursor) break;
  1392. cursor = next;
  1393. }
  1394. let syncedMessages = 0;
  1395. let removedImportedMessages = 0;
  1396. for (const [contactId, group] of grouped.entries()) {
  1397. const conversation = workbench.db.ensureConversation(contactId, group.contactName || '白名单企微联系人');
  1398. removedImportedMessages += workbench.db.deleteImportedMessages(conversation.id, 'manual_sync');
  1399. const ordered = group.messages.sort((a, b) => Date.parse(a.timestamp) - Date.parse(b.timestamp));
  1400. const recent = ordered.slice(-20);
  1401. for (const message of recent) {
  1402. const inserted = workbench.db.insertMessage({
  1403. conversationId: conversation.id,
  1404. externalId: message.externalId,
  1405. direction: message.inbound ? 'inbound' : 'outbound',
  1406. senderType: message.inbound ? 'customer' : 'human',
  1407. content: message.content,
  1408. status: message.inbound ? 'received' : 'sent',
  1409. createdAt: message.timestamp,
  1410. raw: message.raw,
  1411. });
  1412. if (inserted.created) syncedMessages += 1;
  1413. }
  1414. }
  1415. const duplicatesRemoved = workbench.db.cleanupInboundContentDuplicates(60);
  1416. workbench.db.audit({
  1417. actor: 'human',
  1418. action: 'conversation_history_synced',
  1419. detail: { scope: syncScope.scope, contactId: syncScope.contactId || null, conversations: grouped.size, insertedMessages: syncedMessages, removedImportedMessages, scannedMessages, duplicatesRemoved },
  1420. });
  1421. return {
  1422. status: 'ok',
  1423. assistantMessage: grouped.size
  1424. ? syncScope.scope === 'conversation'
  1425. ? '已采集当前个人会话的最近消息;只入库,不运行 Agent、不发送回复'
  1426. : `已补采 ${grouped.size} 个白名单真实会话的最近消息;只入库,不运行 Agent、不发送回复`
  1427. : syncScope.scope === 'conversation'
  1428. ? '当前个人会话暂无可采集的文字消息,请先在企微中收发一条消息'
  1429. : '未读取到白名单联系人的文字会话,请先在企微中与该联系人收发一条消息',
  1430. data: {
  1431. scope: syncScope.scope,
  1432. conversationId: syncScope.conversationId || null,
  1433. syncedConversationCount: grouped.size,
  1434. syncedMessageCount: syncedMessages,
  1435. scannedMessageCount: scannedMessages,
  1436. conversations: sortConversationsByRecency(displayableConversations().map(publicConversation)),
  1437. },
  1438. };
  1439. }
  1440. function changeGlobalMode(mode, confirmation = '') {
  1441. const selected = String(mode || '');
  1442. requireAutopilotConfirmation(selected, confirmation);
  1443. let update;
  1444. if (['paused', 'monitor'].includes(selected)) update = { paused: true };
  1445. else if (['review', 'suggest'].includes(selected)) update = { paused: false, defaultMode: 'review' };
  1446. else if (selected === 'auto') update = { paused: false, defaultMode: 'auto' };
  1447. else if (selected === 'autopilot') update = { paused: false, defaultMode: 'autopilot' };
  1448. else if (selected === 'human') update = { paused: false, defaultMode: 'human' };
  1449. else throw new Error('不支持的 Agent 模式');
  1450. workbench.db.setSetting('listener_enabled', selected === 'paused' ? 'false' : 'true');
  1451. const global = workbench.service.setGlobal(update);
  1452. return {
  1453. status: 'ok',
  1454. assistantMessage: global.paused ? 'Agent 已全局暂停;仍可接收消息,但不会生成或发送回复' : `全局策略已切换为 ${global.defaultMode}`,
  1455. data: { global, globalMode: global.paused ? 'paused' : global.defaultMode },
  1456. };
  1457. }
  1458. function changeConversationMode(id, mode, confirmation = '') {
  1459. const mapped = mode === 'agent' ? 'review' : mode === 'manual' ? 'human' : mode;
  1460. requireAutopilotConfirmation(mapped, confirmation, 'conversation');
  1461. const conversation = workbench.service.setConversationMode(id, mapped);
  1462. const labels = { review: '待审核', auto: '高置信自动', autopilot: '全自动接管', human: '人工接管', paused: '会话暂停' };
  1463. return { status: 'ok', assistantMessage: `会话已切换为${labels[mapped]}`, data: { conversation } };
  1464. }
  1465. async function approveReply(id, content) {
  1466. const detail = workbench.service.conversationDetail(id);
  1467. if (!detail) throw new Error('会话不存在');
  1468. const pending = detail.drafts.find(item => item.status === 'pending');
  1469. const result = pending
  1470. ? await workbench.service.approveDraft(pending.id, { content, actor: 'human' })
  1471. : { status: 'sent', message: await workbench.service.manualSend(id, content, 'human') };
  1472. return { status: 'ok', assistantMessage: '回复已真实发送给白名单测试联系人', data: result };
  1473. }
  1474. async function approveDraft(draftId, content) {
  1475. const result = await workbench.service.approveDraft(draftId, { content, actor: 'human' });
  1476. return { status: 'ok', assistantMessage: '审核通过,回复已真实发送且只发送一次', data: result };
  1477. }
  1478. function rejectDraft(draftId, reason) {
  1479. const result = workbench.service.rejectDraft(draftId, { reason, actor: 'human' });
  1480. return { status: 'ok', assistantMessage: '草稿已驳回,不会发送给客户', data: result };
  1481. }
  1482. async function regenerateDraft(draftId) {
  1483. const result = await workbench.service.regenerateDraft(draftId, 'human');
  1484. return {
  1485. status: 'ok',
  1486. assistantMessage: result.status === 'pending_review' ? 'Agent 已重新生成待审核草稿' : (result.error || 'Agent 重新生成未完成'),
  1487. data: result,
  1488. };
  1489. }
  1490. async function generateLatestDraft(conversationId) {
  1491. const result = await workbench.service.generateLatestDraft(conversationId, 'human');
  1492. return {
  1493. status: 'ok',
  1494. assistantMessage: result.assistantMessage || (result.status === 'pending_review' ? 'Agent 已生成待审核草稿' : (result.error || 'Agent 未生成草稿')),
  1495. data: result,
  1496. };
  1497. }
  1498. async function manualSend(conversationId, content) {
  1499. const message = await workbench.service.manualSend(conversationId, content, 'human');
  1500. return { status: 'ok', assistantMessage: '人工回复已真实发送给白名单测试联系人', data: { message } };
  1501. }
  1502. function getVoiceStatus() {
  1503. return { status: 'ok', data: workbench.voice.status() };
  1504. }
  1505. async function enrollVoice(input = {}) {
  1506. const result = await workbench.voice.enroll({
  1507. filePath: input.filePath,
  1508. originalName: input.originalName,
  1509. mime: input.mime,
  1510. });
  1511. workbench.db.audit({
  1512. actor: 'human',
  1513. action: 'voice_profile_enrolled',
  1514. detail: { duration: result.profile?.duration, sampleRate: result.profile?.sampleRate },
  1515. });
  1516. return {
  1517. status: 'ok',
  1518. assistantMessage: '本人声音已初始化,可以生成企微语音',
  1519. data: workbench.voice.status(),
  1520. };
  1521. }
  1522. function revokeVoiceProfile() {
  1523. const result = workbench.voice.revoke();
  1524. workbench.db.audit({ actor: 'human', action: 'voice_profile_revoked', detail: {} });
  1525. return { status: 'ok', assistantMessage: '声音档案已删除', data: result };
  1526. }
  1527. function voiceConversationContext(conversationId) {
  1528. const conversation = workbench.db.getConversation(conversationId);
  1529. if (!conversation) throw new Error('会话不存在');
  1530. workbench.service.requireAllowed(conversation.contact_id);
  1531. workbench.service.requirePrivateConversation(conversation.id);
  1532. const inbound = workbench.db.getLatestInbound(conversationId);
  1533. return { conversation, context: inbound?.content || '' };
  1534. }
  1535. function pendingVoiceDraft(db, conversationId, draftId = '') {
  1536. const selectedId = String(draftId || '').trim();
  1537. const draft = selectedId
  1538. ? db.getDraft(selectedId)
  1539. : db.listDrafts({ conversationId, status: 'pending', limit: 1 })[0];
  1540. if (!draft && selectedId) throw new Error('待审核草稿不存在或已被处理,请刷新后重试');
  1541. if (!draft) return null;
  1542. if (draft.conversation_id !== conversationId) throw new Error('待审核草稿与当前会话不匹配');
  1543. if (draft.status !== 'pending') throw new Error('待审核草稿已被处理,请刷新后重试');
  1544. return draft;
  1545. }
  1546. function markVoiceDraftSent(db, draft, { content, messageId, actor = 'human:voice' }) {
  1547. if (!draft) return null;
  1548. return db.updateDraft(draft.id, {
  1549. content,
  1550. status: 'sent',
  1551. reviewed_at: new Date().toISOString(),
  1552. reviewer: actor,
  1553. sent_message_id: messageId,
  1554. error: null,
  1555. });
  1556. }
  1557. async function sendClonedVoice(conversationId, input = {}) {
  1558. const { conversation, context } = voiceConversationContext(conversationId);
  1559. const content = String(input.content || '').trim();
  1560. const draft = pendingVoiceDraft(workbench.db, conversationId, input.draftId);
  1561. const result = await workbench.voice.synthesizeAndSend({
  1562. text: content,
  1563. context,
  1564. tone: input.tone || 'auto',
  1565. toId: conversation.contact_id,
  1566. confirmed: input.confirmed,
  1567. });
  1568. if (result.sendResult?.isSendSuccess === false || Number(result.sendResult?.isSendSuccess) === 0) {
  1569. throw new Error('企微未确认语音发送成功');
  1570. }
  1571. const message = workbench.db.insertMessage({
  1572. conversationId,
  1573. direction: 'outbound',
  1574. senderType: 'human',
  1575. content: `[语音 · ${result.tone.label}] ${content}`,
  1576. status: 'sent',
  1577. raw: {
  1578. msgType: 16,
  1579. tone: result.tone,
  1580. duration: result.duration,
  1581. fileSize: result.fileSize,
  1582. audioPath: result.audioPath,
  1583. sendResult: result.sendResult,
  1584. },
  1585. }).message;
  1586. const resolvedDraft = markVoiceDraftSent(workbench.db, draft, { content, messageId: message.id });
  1587. workbench.db.audit({
  1588. actor: 'human',
  1589. action: 'cloned_voice_sent',
  1590. conversationId,
  1591. entityId: message.id,
  1592. detail: {
  1593. tone: result.tone.id,
  1594. automatic: result.tone.automatic,
  1595. duration: result.duration,
  1596. msgServerId: result.sendResult?.msgServerId,
  1597. draftId: resolvedDraft?.id || null,
  1598. },
  1599. });
  1600. const publicMessage = { ...message };
  1601. delete publicMessage.raw_json;
  1602. return {
  1603. status: 'ok',
  1604. assistantMessage: `已用${result.tone.label}语气发送企微语音`,
  1605. data: { message: publicMessage, draft: resolvedDraft, tone: result.tone, duration: result.duration, sendResult: result.sendResult },
  1606. };
  1607. }
  1608. function sentVoiceRuns(voiceRoot = categoryDir('voice')) {
  1609. const runs = [];
  1610. if (!fs.existsSync(voiceRoot)) return runs;
  1611. for (const dateEntry of fs.readdirSync(voiceRoot, { withFileTypes: true })) {
  1612. if (!dateEntry.isDirectory() || !/^\d{4}-\d{2}-\d{2}$/.test(dateEntry.name)) continue;
  1613. const dateDir = path.join(voiceRoot, dateEntry.name);
  1614. for (const runEntry of fs.readdirSync(dateDir, { withFileTypes: true })) {
  1615. if (!runEntry.isDirectory() || !runEntry.name.includes('-clone')) continue;
  1616. const runDir = path.join(dateDir, runEntry.name);
  1617. const manifestPath = path.join(runDir, 'manifest.json');
  1618. const audioPath = path.join(runDir, 'speech.wav');
  1619. if (!fs.existsSync(manifestPath) || !fs.existsSync(audioPath)) continue;
  1620. const manifest = parseJson(fs.readFileSync(manifestPath, 'utf8'), {});
  1621. if (manifest.type !== 'voice-clone' || manifest.outcome === 'failed') continue;
  1622. const createdAt = Date.parse(manifest.createdAt);
  1623. if (!Number.isFinite(createdAt)) continue;
  1624. const textHash = String(manifest.textSha256 || (manifest.text
  1625. ? crypto.createHash('sha256').update(String(manifest.text)).digest('hex')
  1626. : ''));
  1627. if (!textHash) continue;
  1628. runs.push({
  1629. audioPath,
  1630. createdAt,
  1631. duration: Number(manifest.audio?.duration || manifest.duration) || 0,
  1632. textHash,
  1633. });
  1634. }
  1635. }
  1636. return runs;
  1637. }
  1638. function backfillSentVoiceAudioPaths(db, options = {}) {
  1639. const available = sentVoiceRuns(options.voiceRoot).map(item => ({ ...item, used: false }));
  1640. if (!available.length) return 0;
  1641. let updated = 0;
  1642. for (const conversation of db.listConversations()) {
  1643. for (const message of db.listMessages(conversation.id, 1000)) {
  1644. const raw = parseJson(message.raw_json, {});
  1645. if (message.direction !== 'outbound' || message.status !== 'sent'
  1646. || Number(raw.msgType) !== 16 || raw.audioPath) continue;
  1647. const spokenText = String(message.content || '').replace(/^\[语音[^\]]*\]\s*/, '');
  1648. const textHash = crypto.createHash('sha256').update(spokenText).digest('hex');
  1649. const sentAt = Date.parse(message.created_at);
  1650. const match = available
  1651. .filter(item => !item.used && item.textHash === textHash
  1652. && Math.abs(item.createdAt - sentAt) <= 10 * 60 * 1000
  1653. && (!item.duration || !raw.duration || Math.abs(item.duration - Number(raw.duration)) <= 1.5))
  1654. .sort((a, b) => Math.abs(a.createdAt - sentAt) - Math.abs(b.createdAt - sentAt))[0];
  1655. if (!match) continue;
  1656. match.used = true;
  1657. db.updateMessageRaw(message.id, { ...raw, audioPath: match.audioPath });
  1658. updated += 1;
  1659. }
  1660. }
  1661. return updated;
  1662. }
  1663. function getSentVoiceAudio(messageId) {
  1664. const message = workbench.db.getMessage(String(messageId || ''));
  1665. const raw = parseJson(message?.raw_json, {});
  1666. if (!message || message.direction !== 'outbound' || message.status !== 'sent'
  1667. || Number(raw.msgType) !== 16 || !raw.audioPath) {
  1668. throw new Error('已发送语音不存在或不可播放');
  1669. }
  1670. return { filePath: raw.audioPath, duration: Number(raw.duration) || 0 };
  1671. }
  1672. async function updateCustomerTask(taskId, input = {}) {
  1673. if (!['open', 'in_progress', 'done', 'dismissed'].includes(String(input.status || ''))) throw new Error('不支持的客户待办状态');
  1674. const existing = workbench.db.getCustomerTask(taskId);
  1675. if (!existing) throw new Error('客户待办不存在');
  1676. let officialResult = null;
  1677. if (input.status === 'done' && existing.official_todo_id) {
  1678. officialResult = await completeTodoKnowledge(existing.official_todo_id);
  1679. }
  1680. const officialOk = !officialResult || officialResult.status === 'ok';
  1681. const task = workbench.db.updateCustomerTask(taskId, {
  1682. status: input.status,
  1683. resolution_reason: input.status === 'done' ? 'human_completed' : input.status === 'dismissed' ? 'human_dismissed' : '',
  1684. ...(existing.official_todo_id ? { official_sync_status: officialOk ? input.status : 'error' } : {}),
  1685. });
  1686. if (!task) throw new Error('客户待办不存在');
  1687. 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 } });
  1688. return {
  1689. status: 'ok',
  1690. assistantMessage: existing.official_todo_id && !officialOk
  1691. ? `本地待办已更新为 ${task.status},但企微官方待办同步失败,请稍后重试`
  1692. : `客户待办已更新为 ${task.status}${existing.official_todo_id ? ',企微官方待办已同步' : ''}`,
  1693. data: { task, officialResult },
  1694. warnings: existing.official_todo_id && !officialOk ? ['企微官方待办状态尚未同步'] : [],
  1695. };
  1696. }
  1697. async function syncCustomerTaskToOfficialTodo(taskId, input = {}) {
  1698. const syncCurrentAccountTask = createCustomerTaskOfficialSync({
  1699. db: workbench.db,
  1700. searchTodoUsers,
  1701. createTodoKnowledge,
  1702. });
  1703. return syncCurrentAccountTask(taskId, input);
  1704. }
  1705. function updateCustomerAlert(alertId, input = {}) {
  1706. if (!['open', 'acknowledged', 'resolved', 'dismissed'].includes(String(input.status || ''))) throw new Error('不支持的客户预警状态');
  1707. const alert = workbench.db.updateCustomerAlert(alertId, { status: input.status });
  1708. if (!alert) throw new Error('客户预警不存在');
  1709. workbench.db.audit({ actor: 'human', action: 'customer_alert_updated', conversationId: alert.conversation_id, entityId: alert.id, detail: { status: alert.status } });
  1710. return { status: 'ok', assistantMessage: `客户预警已更新为 ${alert.status}`, data: { alert } };
  1711. }
  1712. function addCustomerMemory(conversationId, input = {}) {
  1713. const result = workbench.service.addCustomerMemory(conversationId, input, 'human');
  1714. return { status: 'ok', assistantMessage: '客户长期记忆已添加,并将在后续回复中生效', data: result };
  1715. }
  1716. function updateCustomerMemory(memoryId, input = {}) {
  1717. const action = String(input.action || 'update');
  1718. if (action === 'forget') {
  1719. const result = workbench.service.forgetCustomerMemory(memoryId, 'human');
  1720. return { status: 'ok', assistantMessage: '客户记忆已彻底遗忘', data: result };
  1721. }
  1722. const patch = action === 'confirm'
  1723. ? { type: 'fact', status: 'active', confidence: 1 }
  1724. : action === 'reject'
  1725. ? { status: 'rejected' }
  1726. : {
  1727. content: input.content,
  1728. type: input.type,
  1729. status: input.status,
  1730. importance: input.importance,
  1731. expiresAt: input.expiresAt,
  1732. };
  1733. const sanitized = Object.fromEntries(Object.entries(patch).filter(([, value]) => value !== undefined));
  1734. const result = workbench.service.updateCustomerMemory(memoryId, sanitized, 'human');
  1735. const message = action === 'confirm' ? '推断已人工确认为客户事实' : action === 'reject' ? '客户记忆已拒绝,不再进入上下文' : '客户记忆已更新';
  1736. return { status: 'ok', assistantMessage: message, data: result };
  1737. }
  1738. function getAudit(limit = 200) {
  1739. return { status: 'ok', data: { audit: workbench.db.listAudit(Math.max(1, Math.min(500, Number(limit) || 200))) } };
  1740. }
  1741. async function startListenerForWorkbench(target, account, options = {}) {
  1742. if (!account.online) throw new Error(`${account.nickname || '当前账号'}不在线,无法启动真实消息监听`);
  1743. refreshAllowedSendersFromEnv(target.config.qiwei, options.envFile || ENV_FILE);
  1744. if (options.automatic === true && target.db?.getSetting?.('listener_enabled', 'true') === 'false') {
  1745. return { status: 'ok', assistantMessage: 'AI 监听保持人工关闭状态', data: { running: false, disabled: true } };
  1746. }
  1747. target.db?.setSetting?.('listener_enabled', 'true');
  1748. target.service?.setGlobal?.({ paused: false }, options.automatic ? 'runtime:auto-start' : 'human');
  1749. const status = await target.poller.start();
  1750. return {
  1751. status: 'ok',
  1752. assistantMessage: 'AI 监听已启动:白名单私聊按当前策略处理;已确认客户群自动生成待审核草稿,不会自动群发',
  1753. data: status,
  1754. };
  1755. }
  1756. async function ingestRuntimeMessage(message = {}, options = {}) {
  1757. return ingestMessageForWorkbench({
  1758. config: workbench.config.qiwei,
  1759. db: workbench.db,
  1760. service: workbench.service,
  1761. qiwei: workbench.qiwei,
  1762. }, message, String(options.source || 'runtime'));
  1763. }
  1764. async function startListener(options = {}) {
  1765. const product = getProductMode();
  1766. if (product.mode === 'enterprise') {
  1767. const runtime = listenerStateWithRuntime(product, workbench.poller.status());
  1768. if (!runtime.running) throw new Error('Enterprise Relay runtime is not running. Start the Qiwei runtime first.');
  1769. workbench.service.setGlobal({ paused: false, defaultMode: 'review' });
  1770. return {
  1771. status: 'ok',
  1772. assistantMessage: 'Enterprise Relay keeps collecting messages; AI review and reply processing is enabled.',
  1773. data: runtime,
  1774. };
  1775. }
  1776. const account = await detectAccountStatus(true);
  1777. return startListenerForWorkbench(workbench, account, options);
  1778. }
  1779. function applyManualTakeover(target) {
  1780. target.service.setGlobal({ paused: true, defaultMode: 'review' });
  1781. for (const conversation of target.db.listConversations()) {
  1782. target.service.setConversationMode(conversation.id, 'human');
  1783. }
  1784. }
  1785. function stopListener({ preserveAgentState = false } = {}) {
  1786. const product = getProductMode();
  1787. if (product.mode === 'enterprise') {
  1788. if (!preserveAgentState) {
  1789. workbench.service.setGlobal({ paused: true, defaultMode: 'review' });
  1790. for (const conversation of workbench.db.listConversations()) {
  1791. workbench.service.setConversationMode(conversation.id, 'human');
  1792. }
  1793. }
  1794. return {
  1795. status: 'ok',
  1796. assistantMessage: preserveAgentState
  1797. ? 'Enterprise Relay runtime stopped; Agent modes were preserved.'
  1798. : 'Enterprise Relay continues collecting messages; AI replies are paused and conversations are in human mode.',
  1799. data: { running: false, relayRunning: true, transport: 'server_relay' },
  1800. };
  1801. }
  1802. const status = workbench.poller.stop();
  1803. if (!preserveAgentState) {
  1804. workbench.db.setSetting('listener_enabled', 'false');
  1805. applyManualTakeover(workbench);
  1806. }
  1807. return {
  1808. status: 'ok',
  1809. assistantMessage: preserveAgentState
  1810. ? 'Qiwei polling runtime stopped; Agent modes were preserved.'
  1811. : 'AI 监听已关闭,现有会话已切换为人工接管',
  1812. data: status,
  1813. };
  1814. }
  1815. function getAgentRuntimeConfig() {
  1816. return { ...workbench.config.agent };
  1817. }
  1818. module.exports = {
  1819. switchActiveAccount,
  1820. getAgentStatus,
  1821. getAllowlistCandidates,
  1822. updateAllowlist,
  1823. addAllowlistContacts,
  1824. getIntakePolicy,
  1825. updateIntakePolicy,
  1826. retryOnboardingWelcome,
  1827. getConversations,
  1828. getResponseMonitor,
  1829. updateCustomerProfile,
  1830. syncConversations,
  1831. changeGlobalMode,
  1832. changeConversationMode,
  1833. approveReply,
  1834. approveDraft,
  1835. rejectDraft,
  1836. regenerateDraft,
  1837. generateLatestDraft,
  1838. manualSend,
  1839. getVoiceStatus,
  1840. enrollVoice,
  1841. revokeVoiceProfile,
  1842. sendClonedVoice,
  1843. getSentVoiceAudio,
  1844. updateCustomerTask,
  1845. syncCustomerTaskToOfficialTodo,
  1846. updateCustomerAlert,
  1847. addCustomerMemory,
  1848. updateCustomerMemory,
  1849. getAudit,
  1850. ingestRuntimeMessage,
  1851. ingestWebhookMessage,
  1852. startListener,
  1853. stopListener,
  1854. getAgentRuntimeConfig,
  1855. createWorkbench,
  1856. __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 },
  1857. };