feat: 完成企业微信会话存档全链路功能迭代

1. 后端服务:
   - 新增ArchiveMessage实体与Mapper,实现消息持久化
   - 重构拉取解密逻辑,支持批量解密提升性能
   - 补充员工/客户昵称自动填充与实时同步
   - 修复时区问题,统一使用Asia/Shanghai时区
   - 新增回调线程池,优化异步处理能力
   - 新增批量解密worker命令,解决长参数溢出问题
   - 优化日志打印与异常捕获逻辑

2. 前端页面:
   - 优化聊天时间线展示,显示发送者昵称
   - 重构话术推荐接口调用逻辑,降低接口延迟
   - 优化消息搜索页面,支持昵称搜索与展示
This commit is contained in:
jiao 2026-06-03 14:53:57 +08:00
parent 4d8bc95063
commit 60b6370cd1
34 changed files with 467 additions and 88 deletions

View File

@ -16,6 +16,8 @@ import java.security.MessageDigest;
import java.util.Arrays; import java.util.Arrays;
import java.util.List; import java.util.List;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;
import java.util.concurrent.ExecutorService;
/** /**
* 企微回调控制器 * 企微回调控制器
@ -30,6 +32,12 @@ public class ArchiveCallbackController {
@Autowired @Autowired
private ArchivePullService archivePullService; private ArchivePullService archivePullService;
private final ExecutorService callbackExecutor = Executors.newFixedThreadPool(4, r -> {
Thread t = new Thread(r, "archive-callback-" + System.currentTimeMillis());
t.setDaemon(true);
return t;
});
@Value("${wecom.archive.callback-token:}") @Value("${wecom.archive.callback-token:}")
private String callbackToken; private String callbackToken;
@ -134,7 +142,7 @@ public class ArchiveCallbackController {
} catch (Exception e) { } catch (Exception e) {
log.error("回调触发拉取失败: {}", e.getMessage(), e); log.error("回调触发拉取失败: {}", e.getMessage(), e);
} }
}); }, callbackExecutor);
} }
return Result.success("success"); return Result.success("success");

View File

@ -6,6 +6,7 @@ import com.artedu.archive.entity.Staff;
import com.artedu.archive.mapper.ArchiveMessageMapper; import com.artedu.archive.mapper.ArchiveMessageMapper;
import com.artedu.archive.mapper.CustomerMapper; import com.artedu.archive.mapper.CustomerMapper;
import com.artedu.archive.mapper.StaffMapper; import com.artedu.archive.mapper.StaffMapper;
import com.artedu.archive.client.WeComApiClient;
import com.artedu.archive.service.ArchivePullService; import com.artedu.archive.service.ArchivePullService;
import com.artedu.archive.vo.ArchiveMessageVO; import com.artedu.archive.vo.ArchiveMessageVO;
import com.artedu.common.result.PageResult; import com.artedu.common.result.PageResult;
@ -15,6 +16,7 @@ import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.extension.plugins.pagination.Page; import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.http.HttpHeaders; import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpStatus; import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType; import org.springframework.http.MediaType;
@ -58,6 +60,12 @@ public class ArchiveMessageController {
@Autowired @Autowired
private ArchivePullService archivePullService; private ArchivePullService archivePullService;
@Autowired
private WeComApiClient weComApiClient;
@Value("${wecom.archive.corp-id}")
private String corpId;
private static final String MEDIA_CACHE_DIR = "/app/media"; private static final String MEDIA_CACHE_DIR = "/app/media";
/** /**
@ -158,7 +166,31 @@ public class ArchiveMessageController {
wrapper.eq(ArchiveMessage::getMsgtype, msgtype); wrapper.eq(ArchiveMessage::getMsgtype, msgtype);
} }
if (fromUser != null && !fromUser.isEmpty()) { if (fromUser != null && !fromUser.isEmpty()) {
wrapper.eq(ArchiveMessage::getFromUser, fromUser); // 支持发送人ID精确匹配 + 昵称模糊匹配
Set<String> matchedIds = new java.util.HashSet<>();
matchedIds.add(fromUser);
// 员工昵称模糊匹配
LambdaQueryWrapper<Staff> staffWrapper = new LambdaQueryWrapper<>();
staffWrapper.eq(Staff::getCorpId, corpId).like(Staff::getName, fromUser);
List<Staff> staffs = staffMapper.selectList(staffWrapper);
if (staffs != null) {
for (Staff s : staffs) {
matchedIds.add(s.getStaffId());
}
}
// 客户昵称模糊匹配
LambdaQueryWrapper<Customer> customerWrapper = new LambdaQueryWrapper<>();
customerWrapper.eq(Customer::getCorpId, corpId).like(Customer::getName, fromUser);
List<Customer> customers = customerMapper.selectList(customerWrapper);
if (customers != null) {
for (Customer c : customers) {
matchedIds.add(c.getCustomerId());
}
}
wrapper.in(ArchiveMessage::getFromUser, matchedIds);
} }
if (fromRole != null && !fromRole.isEmpty()) { if (fromRole != null && !fromRole.isEmpty()) {
wrapper.eq(ArchiveMessage::getFromRole, fromRole); wrapper.eq(ArchiveMessage::getFromRole, fromRole);
@ -196,10 +228,31 @@ public class ArchiveMessageController {
private void enrichSenderNames(List<ArchiveMessageVO> list) { private void enrichSenderNames(List<ArchiveMessageVO> list) {
if (list == null || list.isEmpty()) return; if (list == null || list.isEmpty()) return;
Set<String> userIds = list.stream() // 收集所有需要查询昵称的userId(发送者 + 接收者)
.map(ArchiveMessageVO::getFromUser) Set<String> userIds = new java.util.HashSet<>();
.filter(id -> id != null && !id.isEmpty()) for (ArchiveMessageVO vo : list) {
.collect(Collectors.toSet()); // 发送者
String from = vo.getFromUser();
if (from != null && !from.isEmpty()) {
userIds.add(from);
}
// 接收者(解析JSON数组)
String to = vo.getToUser();
if (to != null && !to.isEmpty()) {
if (to.startsWith("[")) {
try {
List<String> toList = JsonUtils.fromJsonList(to, String.class);
if (toList != null) {
userIds.addAll(toList);
}
} catch (Exception e) {
userIds.add(to);
}
} else {
userIds.add(to);
}
}
}
if (userIds.isEmpty()) return; if (userIds.isEmpty()) return;
List<Staff> staffs = staffMapper.selectList( List<Staff> staffs = staffMapper.selectList(
@ -216,12 +269,67 @@ public class ArchiveMessageController {
nameMap.putAll(staffNameMap); nameMap.putAll(staffNameMap);
nameMap.putAll(customerNameMap); nameMap.putAll(customerNameMap);
for (ArchiveMessageVO vo : list) { // 对未找到昵称或昵称为ID的外部客户,实时调用企微API获取并更新数据库
String uid = vo.getFromUser(); for (String uid : userIds) {
if (uid != null && !uid.isEmpty()) { String currentName = nameMap.get(uid);
vo.setFromUserName(nameMap.getOrDefault(uid, uid)); if (isExternalUserId(uid) && (currentName == null || currentName.equals(uid))) {
try {
Map<String, Object> detail = weComApiClient.getExternalContact(uid);
if (detail != null && !detail.isEmpty()) {
Map<String, Object> contact = (Map<String, Object>) detail.get("external_contact");
if (contact != null) {
String name = (String) contact.get("name");
if (name != null && !name.isEmpty()) {
nameMap.put(uid, name);
// 更新数据库
Customer customer = customerMapper.selectOne(
new LambdaQueryWrapper<Customer>()
.eq(Customer::getCustomerId, uid)
.eq(Customer::getCorpId, "wwd483c2fba24ae30a")
);
if (customer != null) {
customer.setName(name);
customerMapper.updateById(customer);
}
}
}
}
} catch (Exception e) {
log.debug("实时获取客户昵称失败: {}", uid);
}
} }
} }
for (ArchiveMessageVO vo : list) {
// 填充发送者昵称
String from = vo.getFromUser();
if (from != null && !from.isEmpty()) {
vo.setFromUserName(nameMap.getOrDefault(from, from));
}
// 填充接收者昵称
String to = vo.getToUser();
if (to != null && !to.isEmpty()) {
if (to.startsWith("[")) {
try {
List<String> toList = JsonUtils.fromJsonList(to, String.class);
if (toList != null && !toList.isEmpty()) {
String toNames = toList.stream()
.map(id -> nameMap.getOrDefault(id, id))
.collect(Collectors.joining(", "));
vo.setToUserName(toNames);
}
} catch (Exception e) {
vo.setToUserName(nameMap.getOrDefault(to, to));
}
} else {
vo.setToUserName(nameMap.getOrDefault(to, to));
}
}
}
}
private boolean isExternalUserId(String userId) {
return userId != null && (userId.startsWith("wm") || userId.startsWith("wo"));
} }
/** /**
@ -284,7 +392,7 @@ public class ArchiveMessageController {
vo.setDecryptStatus(m.getDecryptStatus()); vo.setDecryptStatus(m.getDecryptStatus());
vo.setCreatedAt(m.getCreatedAt() != null ? m.getCreatedAt().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")) : null); vo.setCreatedAt(m.getCreatedAt() != null ? m.getCreatedAt().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")) : null);
if (m.getMsgtime() != null) { if (m.getMsgtime() != null) {
vo.setMsgTimeStr(LocalDateTime.ofInstant(Instant.ofEpochMilli(m.getMsgtime()), ZoneId.systemDefault()) vo.setMsgTimeStr(LocalDateTime.ofInstant(Instant.ofEpochMilli(m.getMsgtime()), ZoneId.of("Asia/Shanghai"))
.format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"))); .format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")));
} }
return vo; return vo;

View File

@ -14,6 +14,7 @@ import org.springframework.web.bind.annotation.RestController;
import java.time.LocalDate; import java.time.LocalDate;
import java.time.LocalDateTime; import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter; import java.time.format.DateTimeFormatter;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.HashMap; import java.util.HashMap;
@ -52,7 +53,7 @@ public class DashboardController {
DateTimeFormatter fmt = DateTimeFormatter.ofPattern("MM-dd"); DateTimeFormatter fmt = DateTimeFormatter.ofPattern("MM-dd");
for (int i = days - 1; i >= 0; i--) { for (int i = days - 1; i >= 0; i--) {
LocalDate date = LocalDate.now().minusDays(i); LocalDate date = LocalDate.now(ZoneId.of("Asia/Shanghai")).minusDays(i);
LocalDateTime start = date.atStartOfDay(); LocalDateTime start = date.atStartOfDay();
LocalDateTime end = date.plusDays(1).atStartOfDay(); LocalDateTime end = date.plusDays(1).atStartOfDay();

View File

@ -268,7 +268,8 @@ public class ArchiveContactService {
if (gender != null) customer.setGender(gender == 1 ? "MALE" : gender == 2 ? "FEMALE" : null); if (gender != null) customer.setGender(gender == 1 ? "MALE" : gender == 2 ? "FEMALE" : null);
if (unionid != null) customer.setUnionid(unionid); if (unionid != null) customer.setUnionid(unionid);
if (addTime != null) { if (addTime != null) {
customer.setAddTime(new java.sql.Timestamp(addTime).toLocalDateTime()); customer.setAddTime(java.time.LocalDateTime.ofInstant(
java.time.Instant.ofEpochMilli(addTime), java.time.ZoneId.of("Asia/Shanghai")));
} }
customerMapper.updateById(customer); customerMapper.updateById(customer);
log.info("补充客户资料成功: customerId={}, name={}", customerId, name); log.info("补充客户资料成功: customerId={}, name={}", customerId, name);

View File

@ -17,6 +17,8 @@ import javax.annotation.PreDestroy;
import java.io.File; import java.io.File;
import java.nio.charset.StandardCharsets; import java.nio.charset.StandardCharsets;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
@ -118,7 +120,7 @@ public class ArchivePullService {
Map<String, Object> result = JsonUtils.fromJsonMap(jsonData); Map<String, Object> result = JsonUtils.fromJsonMap(jsonData);
if (result == null || !Integer.valueOf(0).equals(result.get("errcode"))) { if (result == null || !Integer.valueOf(0).equals(result.get("errcode"))) {
log.error("拉取存档返回错误: {}", jsonData); log.error("拉取存档返回错误: errcode={}, errmsg={}, json前200字={}", result != null ? result.get("errcode") : "null", result != null ? result.get("errmsg") : "null", jsonData.length() > 200 ? jsonData.substring(0, 200) : jsonData);
break; break;
} }
@ -127,21 +129,25 @@ public class ArchivePullService {
log.info("已无更多消息, seq={}", currentSeq); log.info("已无更多消息, seq={}", currentSeq);
break; break;
} }
log.info("拉取到消息条数: seq={}, count={}", currentSeq, chatDataList.size());
// 批量解密(一次子进程解密所有消息,避免重复启动 JVM)
List<ArchiveMessage> batchMessages = batchDecryptAndParse(chatDataList);
if (batchMessages != null && !batchMessages.isEmpty()) {
messages.addAll(batchMessages);
totalParsed += batchMessages.size();
}
totalFetched += chatDataList.size();
// 计算 maxSeq
long maxSeq = currentSeq; long maxSeq = currentSeq;
for (Map<String, Object> chatData : chatDataList) { for (Map<String, Object> chatData : chatDataList) {
try { try {
Long msgSeq = ((Number) chatData.get("seq")).longValue(); Object seqObj = chatData.get("seq");
Long msgSeq = seqObj instanceof Number ? ((Number) seqObj).longValue() : Long.parseLong(String.valueOf(seqObj));
if (msgSeq > maxSeq) maxSeq = msgSeq; if (msgSeq > maxSeq) maxSeq = msgSeq;
totalFetched++;
ArchiveMessage message = parseAndDecrypt(chatData);
if (message != null) {
messages.add(message);
totalParsed++;
}
} catch (Exception e) { } catch (Exception e) {
log.error("解析单条消息失败: {}", e.getMessage(), e); log.error("计算seq失败: {}", e.getMessage());
} }
} }
@ -168,25 +174,102 @@ public class ArchivePullService {
Long msgSeq = ((Number) chatData.get("seq")).longValue(); Long msgSeq = ((Number) chatData.get("seq")).longValue();
Object pubKeyVer = chatData.get("publickey_ver"); Object pubKeyVer = chatData.get("publickey_ver");
// publickey_ver 与当前私钥版本不匹配时跳过(企微后台可能更换过公钥)
// 当前私钥对应版本3,历史消息可能是版本2加密的
log.debug("chatData: seq={}, publickey_ver={}", msgSeq, pubKeyVer); log.debug("chatData: seq={}, publickey_ver={}", msgSeq, pubKeyVer);
// 主服务中 RSA 解密(不涉及 JNI,安全)
String encryptKey = rsaDecryptUtil.decrypt(encryptRandomKey, rsaPrivateKey); String encryptKey = rsaDecryptUtil.decrypt(encryptRandomKey, rsaPrivateKey);
if (encryptKey == null) { if (encryptKey == null) {
log.debug("解密encrypt_random_key失败, seq={}, 可能是公钥版本不匹配", msgSeq); log.debug("解密encrypt_random_key失败, seq={}, 可能是公钥版本不匹配", msgSeq);
return null; return null;
} }
// Worker 2: 解密单条消息内容(隔离崩溃风险)
String decryptedMsg = callWorker("decrypt", corpId, secret, sdkPath, encryptKey, encryptChatMsg); String decryptedMsg = callWorker("decrypt", corpId, secret, sdkPath, encryptKey, encryptChatMsg);
if (decryptedMsg == null) { if (decryptedMsg == null) {
log.error("解密消息内容失败(worker崩溃或超时), seq={}", msgSeq); log.error("解密消息内容失败(worker崩溃或超时), seq={}", msgSeq);
return null; return null;
} }
Map<String, Object> msgMap = JsonUtils.fromJsonMap(decryptedMsg); return buildMessageFromJson(chatData, decryptedMsg, msgSeq);
}
/**
* 批量解密:一次子进程解密所有消息,避免重复启动 JVM
*/
private List<ArchiveMessage> batchDecryptAndParse(List<Map<String, Object>> chatDataList) {
List<Map<String, Object>> inputs = new ArrayList<>();
Map<Long, Map<String, Object>> seqToChatData = new HashMap<>();
for (Map<String, Object> chatData : chatDataList) {
Object seqObj = chatData.get("seq");
Long msgSeq = seqObj instanceof Number ? ((Number) seqObj).longValue() : Long.parseLong(String.valueOf(seqObj));
String encryptRandomKey = (String) chatData.get("encrypt_random_key");
String encryptChatMsg = (String) chatData.get("encrypt_chat_msg");
String encryptKey = rsaDecryptUtil.decrypt(encryptRandomKey, rsaPrivateKey);
if (encryptKey == null) {
log.debug("批量解密: RSA解密失败, seq={}", msgSeq);
continue;
}
Map<String, Object> input = new HashMap<>();
input.put("seq", msgSeq);
input.put("encryptKey", encryptKey);
input.put("encryptMsg", encryptChatMsg);
inputs.add(input);
seqToChatData.put(msgSeq, chatData);
}
if (inputs.isEmpty()) {
return Collections.emptyList();
}
String json = JsonUtils.toJson(inputs);
String result = callWorkerWithStdin(new String[]{"batchdecrypt", corpId, secret, sdkPath}, json);
if (result == null) {
log.error("批量解密返回null, count={}", inputs.size());
return Collections.emptyList();
}
List<Map<String, Object>> outputs;
try {
@SuppressWarnings("unchecked")
List<Map<String, Object>> temp = (List<Map<String, Object>>) (List<?>) JsonUtils.fromJsonList(result, Map.class);
outputs = temp;
} catch (Exception e) {
log.error("批量解密结果解析失败: {}", e.getMessage());
return Collections.emptyList();
}
if (outputs == null) {
return Collections.emptyList();
}
List<ArchiveMessage> messages = new ArrayList<>();
for (Map<String, Object> output : outputs) {
Object seqObj = output.get("seq");
Long msgSeq = seqObj instanceof Number ? ((Number) seqObj).longValue() : Long.parseLong(String.valueOf(seqObj));
Object errcodeObj = output.get("errcode");
int errcode = errcodeObj instanceof Number ? ((Number) errcodeObj).intValue() : Integer.parseInt(String.valueOf(errcodeObj));
if (errcode != 0) {
log.error("批量解密单条失败: seq={}, errcode={}", msgSeq, errcode);
continue;
}
String decryptedJson = (String) output.get("data");
Map<String, Object> chatData = seqToChatData.get(msgSeq);
if (chatData == null || decryptedJson == null) continue;
ArchiveMessage message = buildMessageFromJson(chatData, decryptedJson, msgSeq);
if (message != null) {
messages.add(message);
}
}
return messages;
}
private ArchiveMessage buildMessageFromJson(Map<String, Object> chatData, String decryptedJson, Long msgSeq) {
Map<String, Object> msgMap = JsonUtils.fromJsonMap(decryptedJson);
if (msgMap == null) { if (msgMap == null) {
return null; return null;
} }
@ -199,7 +282,6 @@ public class ArchivePullService {
message.setFromUser((String) msgMap.get("from")); message.setFromUser((String) msgMap.get("from"));
message.setFromRole(detectRole((String) msgMap.get("from"))); message.setFromRole(detectRole((String) msgMap.get("from")));
// tolist 可能是字符串或数组
Object tolistObj = msgMap.get("tolist"); Object tolistObj = msgMap.get("tolist");
if (tolistObj instanceof java.util.List) { if (tolistObj instanceof java.util.List) {
message.setToUser(JsonUtils.toJson(tolistObj)); message.setToUser(JsonUtils.toJson(tolistObj));
@ -214,8 +296,6 @@ public class ArchivePullService {
message.setMsgtime(((Number) msgTimeObj).longValue()); message.setMsgtime(((Number) msgTimeObj).longValue());
} }
// 企微消息结构: { msgtype: "text", text: { content: "..." } }
// 根据 msgtype 从对应子对象提取内容
String msgtype = message.getMsgtype(); String msgtype = message.getMsgtype();
Object typeContent = msgMap.get(msgtype); Object typeContent = msgMap.get(msgtype);
if (typeContent instanceof Map) { if (typeContent instanceof Map) {
@ -266,6 +346,10 @@ public class ArchivePullService {
private String callWorker(String... workerArgs) { private String callWorker(String... workerArgs) {
return callWorkerWithStdin(workerArgs, null);
}
private String callWorkerWithStdin(String[] workerArgs, String stdinData) {
List<String> cmd = new ArrayList<>(); List<String> cmd = new ArrayList<>();
cmd.add("java"); cmd.add("java");
cmd.add("-Djava.library.path=" + sdkPath); cmd.add("-Djava.library.path=" + sdkPath);
@ -274,26 +358,48 @@ public class ArchivePullService {
cmd.add("worker"); cmd.add("worker");
cmd.addAll(java.util.Arrays.asList(workerArgs)); cmd.addAll(java.util.Arrays.asList(workerArgs));
Process process = null;
java.util.concurrent.ExecutorService readerExecutor = null;
try { try {
ProcessBuilder pb = new ProcessBuilder(cmd); ProcessBuilder pb = new ProcessBuilder(cmd);
pb.redirectErrorStream(true); pb.redirectErrorStream(true);
Process process = pb.start(); process = pb.start();
String output; // 如果有 stdin 数据,写入子进程标准输入
try (java.io.BufferedReader reader = new java.io.BufferedReader( if (stdinData != null && !stdinData.isEmpty()) {
new java.io.InputStreamReader(process.getInputStream(), StandardCharsets.UTF_8))) { try (java.io.OutputStream os = process.getOutputStream()) {
StringBuilder sb = new StringBuilder(); os.write(stdinData.getBytes(StandardCharsets.UTF_8));
String line; os.flush();
while ((line = reader.readLine()) != null) {
sb.append(line);
} }
output = sb.toString();
} }
boolean finished = process.waitFor(30, TimeUnit.SECONDS); final Process p = process;
readerExecutor = java.util.concurrent.Executors.newSingleThreadExecutor();
java.util.concurrent.Future<String> future = readerExecutor.submit(() -> {
try (java.io.BufferedReader reader = new java.io.BufferedReader(
new java.io.InputStreamReader(p.getInputStream(), StandardCharsets.UTF_8))) {
StringBuilder sb = new StringBuilder();
String line;
while ((line = reader.readLine()) != null) {
sb.append(line);
}
return sb.toString();
}
});
String output;
try {
output = future.get(30, TimeUnit.SECONDS);
} catch (java.util.concurrent.TimeoutException e) {
future.cancel(true);
log.error("Worker 读取超时: {}", java.util.Arrays.toString(workerArgs));
return null;
}
boolean finished = process.waitFor(5, TimeUnit.SECONDS);
if (!finished) { if (!finished) {
process.destroyForcibly(); process.destroyForcibly();
log.error("Worker 超时: {}", java.util.Arrays.toString(workerArgs)); log.error("Worker 进程未在5秒内退出,已强制终止: {}", java.util.Arrays.toString(workerArgs));
return null; return null;
} }
@ -303,10 +409,18 @@ public class ArchivePullService {
return null; return null;
} }
log.info("Worker 执行成功: {}, output长度={}", java.util.Arrays.toString(workerArgs), output != null ? output.length() : 0);
return output; return output;
} catch (Exception e) { } catch (Exception e) {
log.error("Worker 执行失败: {}", java.util.Arrays.toString(workerArgs), e); log.error("Worker 执行失败: {}", java.util.Arrays.toString(workerArgs), e);
return null; return null;
} finally {
if (process != null && process.isAlive()) {
process.destroyForcibly();
}
if (readerExecutor != null) {
readerExecutor.shutdownNow();
}
} }
} }
@ -352,6 +466,8 @@ public class ArchivePullService {
} }
public void saveAndNotify(List<ArchiveMessage> messages) { public void saveAndNotify(List<ArchiveMessage> messages) {
log.info("开始保存消息, 共{}条", messages.size());
int saved = 0;
for (ArchiveMessage message : messages) { for (ArchiveMessage message : messages) {
try { try {
int count = archiveMessageMapper.countByMsgId(message.getMsgid(), message.getCorpId()); int count = archiveMessageMapper.countByMsgId(message.getMsgid(), message.getCorpId());
@ -370,6 +486,7 @@ public class ArchivePullService {
} }
archiveMessageMapper.insert(message); archiveMessageMapper.insert(message);
saved++;
// 提取并维护联系人信息 // 提取并维护联系人信息
archiveContactService.processMessage(message); archiveContactService.processMessage(message);
@ -390,17 +507,25 @@ public class ArchivePullService {
); );
} catch (Exception e) { } catch (Exception e) {
log.error("保存消息失败: msgid={}, error={}", message.getMsgid(), e.getMessage()); log.error("保存消息失败: msgid={}, error={}", message.getMsgid(), e.getMessage(), e);
} }
} }
log.info("保存消息完成, 成功{}条/共{}条", saved, messages.size());
} }
public long getLastSeq() { public long getLastSeq() {
Object seq = redisTemplate.opsForValue().get(LAST_SEQ_KEY + corpId); Object seq = redisTemplate.opsForValue().get(LAST_SEQ_KEY + corpId);
if (seq != null) {
return ((Number) seq).longValue();
}
Long maxSeq = archiveMessageMapper.selectMaxSeq(corpId); Long maxSeq = archiveMessageMapper.selectMaxSeq(corpId);
if (seq != null) {
long redisSeq = ((Number) seq).longValue();
// 防御:如果 Redis 中的 seq 小于数据库最大值,使用数据库最大值
if (maxSeq != null && redisSeq < maxSeq) {
log.warn("Redis seq({}) 小于数据库最大 seq({}),使用数据库最大值", redisSeq, maxSeq);
redisTemplate.opsForValue().set(LAST_SEQ_KEY + corpId, maxSeq, 7, TimeUnit.DAYS);
return maxSeq;
}
return redisSeq;
}
return maxSeq != null ? maxSeq : 0L; return maxSeq != null ? maxSeq : 0L;
} }
@ -446,9 +571,8 @@ public class ArchivePullService {
long startSeq = minSeq - 1; long startSeq = minSeq - 1;
log.info("开始修复历史数据,起始 seq={}", startSeq); log.info("开始修复历史数据,起始 seq={}", startSeq);
// 重置 Redis 中的 last_seq // 修复任务使用自己的 currentSeq,不重置 Redis 中的 last_seq
redisTemplate.opsForValue().set(LAST_SEQ_KEY + corpId, startSeq, 7, TimeUnit.DAYS); // 避免影响正常的增量拉取
int totalFetched = 0; int totalFetched = 0;
int totalUpdated = 0; int totalUpdated = 0;
int totalInserted = 0; int totalInserted = 0;

View File

@ -18,10 +18,10 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.text.SimpleDateFormat;
import java.time.Instant; import java.time.Instant;
import java.time.LocalDateTime; import java.time.LocalDateTime;
import java.time.ZoneId; import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
import java.util.*; import java.util.*;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ -44,7 +44,8 @@ public class ArchiveSessionService {
@Autowired @Autowired
private ArchiveMessageMapper archiveMessageMapper; private ArchiveMessageMapper archiveMessageMapper;
private static final SimpleDateFormat TIME_FORMAT = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); private static final java.time.format.DateTimeFormatter TIME_FORMATTER =
java.time.format.DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
/** /**
* 员工列表 * 员工列表
@ -349,7 +350,8 @@ public class ArchiveSessionService {
private String formatTime(Long msgTime) { private String formatTime(Long msgTime) {
if (msgTime == null) return "-"; if (msgTime == null) return "-";
try { try {
return TIME_FORMAT.format(new Date(msgTime)); return LocalDateTime.ofInstant(Instant.ofEpochMilli(msgTime), ZoneId.of("Asia/Shanghai"))
.format(TIME_FORMATTER);
} catch (Exception e) { } catch (Exception e) {
return String.valueOf(msgTime); return String.valueOf(msgTime);
} }

View File

@ -17,6 +17,7 @@ public class ArchiveMessageVO {
private String fromUserName; private String fromUserName;
private String fromRole; private String fromRole;
private String toUser; private String toUser;
private String toUserName;
private String tolist; private String tolist;
private String roomid; private String roomid;
private String msgtype; private String msgtype;

View File

@ -17,7 +17,7 @@ public class SdkWorker {
public static void run(String[] args) throws Exception { public static void run(String[] args) throws Exception {
if (args.length < 1) { if (args.length < 1) {
System.err.println("Usage: worker <pull|decrypt|media> ..."); System.err.println("Usage: worker <pull|decrypt|batchdecrypt|media> ...");
System.exit(1); System.exit(1);
} }
String cmd = args[0]; String cmd = args[0];
@ -25,6 +25,8 @@ public class SdkWorker {
runPull(args); runPull(args);
} else if ("decrypt".equals(cmd)) { } else if ("decrypt".equals(cmd)) {
runDecrypt(args); runDecrypt(args);
} else if ("batchdecrypt".equals(cmd)) {
runBatchDecrypt(args);
} else if ("media".equals(cmd)) { } else if ("media".equals(cmd)) {
runGetMediaData(args); runGetMediaData(args);
} else { } else {
@ -89,6 +91,67 @@ public class SdkWorker {
} }
} }
private static void runBatchDecrypt(String[] args) throws Exception {
if (args.length < 4) {
System.err.println("Usage: worker batchdecrypt <corpId> <secret> <sdkPath>");
System.exit(1);
}
String corpId = args[1];
String secret = args[2];
String sdkPath = args[3];
// 从 stdin 读取 JSON 输入,避免命令行参数过长(error=7 Argument list too long)
StringBuilder sb = new StringBuilder();
byte[] buffer = new byte[8192];
int read;
while ((read = System.in.read(buffer)) != -1) {
sb.append(new String(buffer, 0, read, java.nio.charset.StandardCharsets.UTF_8));
}
String json = sb.toString().trim();
if (json.isEmpty()) {
System.err.println("ERROR batch input is empty");
System.exit(1);
}
setupLibraryPath(sdkPath);
long sdk = initSdk(corpId, secret);
@SuppressWarnings("unchecked")
java.util.List<java.util.Map<String, Object>> inputs = (java.util.List<java.util.Map<String, Object>>) (java.util.List<?>) JsonUtils.fromJsonList(json, java.util.Map.class);
java.util.List<java.util.Map<String, Object>> outputs = new java.util.ArrayList<>();
for (java.util.Map<String, Object> input : inputs) {
String encryptKey = (String) input.get("encryptKey");
String encryptMsg = (String) input.get("encryptMsg");
Object seqObj = input.get("seq");
long seq = seqObj instanceof Number ? ((Number) seqObj).longValue() : Long.parseLong(String.valueOf(seqObj));
long slice = Finance.NewSlice();
try {
int ret = Finance.DecryptData(sdk, encryptKey, encryptMsg, slice);
if (ret != 0) {
java.util.Map<String, Object> err = new java.util.HashMap<>();
err.put("seq", seq);
err.put("errcode", ret);
err.put("errmsg", "DecryptData failed");
outputs.add(err);
continue;
}
String decryptedJson = Finance.GetContentFromSlice(slice);
java.util.Map<String, Object> ok = new java.util.HashMap<>();
ok.put("seq", seq);
ok.put("errcode", 0);
ok.put("data", decryptedJson);
outputs.add(ok);
} finally {
Finance.FreeSlice(slice);
}
}
Finance.DestroySdk(sdk);
System.out.println(JsonUtils.toJson(outputs));
}
private static void runGetMediaData(String[] args) throws Exception { private static void runGetMediaData(String[] args) throws Exception {
if (args.length < 6) { if (args.length < 6) {
System.err.println("Usage: worker media <corpId> <secret> <sdkPath> <sdkfileid> <outputPath>"); System.err.println("Usage: worker media <corpId> <secret> <sdkPath> <sdkfileid> <outputPath>");

View File

@ -1,11 +1,11 @@
package com.artedu.conversation.controller; package com.artedu.conversation.controller;
import com.artedu.common.result.Result; import com.artedu.common.result.Result;
import com.artedu.conversation.entity.ArchiveMessage;
import com.artedu.conversation.entity.Conversation; import com.artedu.conversation.entity.Conversation;
import com.artedu.conversation.entity.ConversationTurn; import com.artedu.conversation.entity.ConversationTurn;
import com.artedu.conversation.mapper.ArchiveMessageMapper;
import com.artedu.conversation.service.ConversationManager; import com.artedu.conversation.service.ConversationManager;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.toolkit.Wrappers;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.validation.annotation.Validated; import org.springframework.validation.annotation.Validated;
@ -25,6 +25,9 @@ public class ConversationController {
@Autowired @Autowired
private ConversationManager conversationManager; private ConversationManager conversationManager;
@Autowired
private ArchiveMessageMapper archiveMessageMapper;
@GetMapping("/{sessionId}/context") @GetMapping("/{sessionId}/context")
public Result<Map<String, Object>> getContext(@NotBlank(message = "sessionId不能为空") @PathVariable("sessionId") String sessionId) { public Result<Map<String, Object>> getContext(@NotBlank(message = "sessionId不能为空") @PathVariable("sessionId") String sessionId) {
String contextStr = conversationManager.getContextString(sessionId); String contextStr = conversationManager.getContextString(sessionId);
@ -59,6 +62,22 @@ public class ConversationController {
return Result.success(list); return Result.success(list);
} }
@GetMapping("/last-message")
public Result<Map<String, Object>> getLastMessage(
@NotBlank(message = "customerId不能为空") @RequestParam("customerId") String customerId,
@NotBlank(message = "corpId不能为空") @RequestParam("corpId") String corpId) {
ArchiveMessage message = archiveMessageMapper.selectLastMessage(customerId, corpId);
Map<String, Object> result = new HashMap<>();
if (message != null) {
result.put("content", message.getContent());
result.put("msgTime", message.getMsgtime());
result.put("msgType", message.getMsgtype());
result.put("fromUser", message.getFromUser());
}
return Result.success(result);
}
@PutMapping("/{sessionId}/status") @PutMapping("/{sessionId}/status")
public Result<String> updateStatus( public Result<String> updateStatus(
@NotBlank(message = "sessionId不能为空") @PathVariable("sessionId") String sessionId, @NotBlank(message = "sessionId不能为空") @PathVariable("sessionId") String sessionId,

View File

@ -0,0 +1,32 @@
package com.artedu.conversation.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.time.LocalDateTime;
@Data
@TableName("archive_messages")
public class ArchiveMessage {
@TableId(type = IdType.AUTO)
private Long id;
private String msgid;
private Long seq;
private String corpId;
private String action;
private String fromUser;
private String fromRole;
private String toUser;
private String tolist;
private String roomid;
private String msgtype;
private Long msgtime;
private String content;
private String mediaData;
private Integer decryptStatus;
private String sessionId;
private LocalDateTime createdAt;
private LocalDateTime updatedAt;
}

View File

@ -0,0 +1,19 @@
package com.artedu.conversation.mapper;
import com.artedu.conversation.entity.ArchiveMessage;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import org.apache.ibatis.annotations.Select;
@Mapper
public interface ArchiveMessageMapper extends BaseMapper<ArchiveMessage> {
@Select("SELECT * FROM archive_messages " +
"WHERE (from_user = #{userId} OR to_user LIKE CONCAT('%', #{userId}, '%')) " +
"AND corp_id = #{corpId} " +
"AND content IS NOT NULL AND content != '' " +
"ORDER BY msgtime DESC LIMIT 1")
ArchiveMessage selectLastMessage(@Param("userId") String userId,
@Param("corpId") String corpId);
}

View File

@ -3,7 +3,9 @@ com\artedu\conversation\ConversationServiceApplication.class
com\artedu\conversation\mapper\ConversationTurnMapper.class com\artedu\conversation\mapper\ConversationTurnMapper.class
com\artedu\conversation\entity\Conversation.class com\artedu\conversation\entity\Conversation.class
com\artedu\conversation\entity\ConversationTurn.class com\artedu\conversation\entity\ConversationTurn.class
com\artedu\conversation\mapper\ArchiveMessageMapper.class
com\artedu\conversation\service\ConversationManager$ContextMessage.class com\artedu\conversation\service\ConversationManager$ContextMessage.class
com\artedu\conversation\entity\ArchiveMessage.class
com\artedu\conversation\service\ConversationManager.class com\artedu\conversation\service\ConversationManager.class
com\artedu\conversation\service\ConversationManager$ArchiveMessageEvent.class com\artedu\conversation\service\ConversationManager$ArchiveMessageEvent.class
com\artedu\conversation\controller\ConversationController.class com\artedu\conversation\controller\ConversationController.class

View File

@ -1,7 +1,9 @@
D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\controller\ConversationController.java D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\controller\ConversationController.java
D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\ConversationServiceApplication.java D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\ConversationServiceApplication.java
D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\entity\ArchiveMessage.java
D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\entity\Conversation.java D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\entity\Conversation.java
D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\entity\ConversationTurn.java D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\entity\ConversationTurn.java
D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\mapper\ArchiveMessageMapper.java
D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\mapper\ConversationMapper.java D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\mapper\ConversationMapper.java
D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\mapper\ConversationTurnMapper.java D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\mapper\ConversationTurnMapper.java
D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\service\ConversationManager.java D:\www\agent_9art\backend\conversation-service\src\main\java\com\artedu\conversation\service\ConversationManager.java

View File

@ -246,6 +246,7 @@ export default function ChatTimeline({ sessionId, staffName = '员工', customer
{renderContent(msg)} {renderContent(msg)}
</div> </div>
<div style={{ fontSize: 11, color: '#aaa', marginTop: 4, textAlign: staffSide ? 'right' : 'left' }}> <div style={{ fontSize: 11, color: '#aaa', marginTop: 4, textAlign: staffSide ? 'right' : 'left' }}>
<span style={{ marginRight: 4, color: '#666' }}>{staffSide ? staffName : customerName}</span>
<Tag style={{ fontSize: 10, padding: '0 4px', lineHeight: '16px' }}> <Tag style={{ fontSize: 10, padding: '0 4px', lineHeight: '16px' }}>
{msgTypeMap[msg.msgtype]?.label || msg.msgtype} {msgTypeMap[msg.msgtype]?.label || msg.msgtype}
</Tag> </Tag>

View File

@ -9,11 +9,14 @@ interface ArchiveMessage {
id: number id: number
msgid: string msgid: string
fromUser: string fromUser: string
fromUserName?: string
fromRole: string fromRole: string
toUser: string toUser: string
toUserName?: string
msgtype: string msgtype: string
content: string content: string
msgtime: number msgtime: number
msgTimeStr?: string
sessionId: string sessionId: string
} }
@ -40,6 +43,8 @@ export default function MessageSearch() {
const [searched, setSearched] = useState(false) const [searched, setSearched] = useState(false)
const [chatVisible, setChatVisible] = useState(false) const [chatVisible, setChatVisible] = useState(false)
const [chatSessionId, setChatSessionId] = useState<string>('') const [chatSessionId, setChatSessionId] = useState<string>('')
const [chatStaffName, setChatStaffName] = useState<string>('')
const [chatCustomerName, setChatCustomerName] = useState<string>('')
const loadMessages = async (page = 1, pageSize = 20, searchParams?: any) => { const loadMessages = async (page = 1, pageSize = 20, searchParams?: any) => {
setLoading(true) setLoading(true)
@ -97,7 +102,8 @@ export default function MessageSearch() {
return dayjs(ts).format('YYYY-MM-DD HH:mm:ss') return dayjs(ts).format('YYYY-MM-DD HH:mm:ss')
} }
const getSenderLabel = (role: string, userId: string) => { const getSenderLabel = (role: string, userId: string, userName?: string) => {
if (userName && userName !== userId) return userName
if (role === 'INTERNAL') return `员工:${userId}` if (role === 'INTERNAL') return `员工:${userId}`
if (role === 'EXTERNAL') return `客户:${userId}` if (role === 'EXTERNAL') return `客户:${userId}`
return userId return userId
@ -129,7 +135,7 @@ export default function MessageSearch() {
width: 150, width: 150,
render: (fromUser: string, record: ArchiveMessage) => ( render: (fromUser: string, record: ArchiveMessage) => (
<Tag color={record.fromRole === 'INTERNAL' ? 'blue' : 'green'}> <Tag color={record.fromRole === 'INTERNAL' ? 'blue' : 'green'}>
{getSenderLabel(record.fromRole, fromUser)} {getSenderLabel(record.fromRole, fromUser, record.fromUserName)}
</Tag> </Tag>
), ),
}, },
@ -138,14 +144,18 @@ export default function MessageSearch() {
dataIndex: 'toUser', dataIndex: 'toUser',
key: 'toUser', key: 'toUser',
width: 150, width: 150,
render: (toUser: string) => toUser || '-', render: (toUser: string, record: ArchiveMessage) => (
record.toUserName || toUser || '-'
),
}, },
{ {
title: '发送时间', title: '发送时间',
dataIndex: 'msgtime', dataIndex: 'msgtime',
key: 'msgtime', key: 'msgtime',
width: 170, width: 170,
render: (msgtime: number) => formatTime(msgtime), render: (msgtime: number, record: ArchiveMessage) => (
record.msgTimeStr || formatTime(msgtime)
),
}, },
{ {
title: '操作', title: '操作',
@ -157,6 +167,14 @@ export default function MessageSearch() {
icon={<MessageOutlined />} icon={<MessageOutlined />}
onClick={() => { onClick={() => {
setChatSessionId(record.sessionId) setChatSessionId(record.sessionId)
// 根据 fromRole 判断谁是员工、谁是客户,避免昵称对调
if (record.fromRole === 'INTERNAL') {
setChatStaffName(record.fromUserName || record.fromUser)
setChatCustomerName(record.toUserName || record.toUser)
} else {
setChatStaffName(record.toUserName || record.toUser)
setChatCustomerName(record.fromUserName || record.fromUser)
}
setChatVisible(true) setChatVisible(true)
}} }}
> >
@ -194,12 +212,12 @@ export default function MessageSearch() {
</Col> </Col>
<Col span={4}> <Col span={4}>
<Form.Item name="fromUser" label="发送者"> <Form.Item name="fromUser" label="发送者">
<Input placeholder="发送者ID" allowClear /> <Input placeholder="发送者昵称/ID" allowClear />
</Form.Item> </Form.Item>
</Col> </Col>
<Col span={4}> <Col span={4}>
<Form.Item name="toUser" label="接收者"> <Form.Item name="toUser" label="接收者">
<Input placeholder="接收者ID" allowClear /> <Input placeholder="接收者昵称/ID" allowClear />
</Form.Item> </Form.Item>
</Col> </Col>
<Col span={4}> <Col span={4}>
@ -251,7 +269,7 @@ export default function MessageSearch() {
footer={null} footer={null}
destroyOnClose destroyOnClose
> >
<ChatTimeline sessionId={chatSessionId} /> <ChatTimeline sessionId={chatSessionId} staffName={chatStaffName} customerName={chatCustomerName} />
</Modal> </Modal>
</div> </div>
) )

View File

@ -7,12 +7,6 @@ interface Props {
customerId?: string customerId?: string
} }
interface ConversationTurn {
turnNumber: number
studentContent: string
seatContent: string
}
export default function ScriptRecommend({ userInfo, customerId: propCustomerId }: Props) { export default function ScriptRecommend({ userInfo, customerId: propCustomerId }: Props) {
const [scripts, setScripts] = useState<Recommendation[]>([]) const [scripts, setScripts] = useState<Recommendation[]>([])
const [loading, setLoading] = useState(false) const [loading, setLoading] = useState(false)
@ -26,30 +20,14 @@ export default function ScriptRecommend({ userInfo, customerId: propCustomerId }
const staffId = userInfo?.userId || 'staff_001' const staffId = userInfo?.userId || 'staff_001'
const corpId = userInfo?.corpId || 'wwd483c2fba24ae30a' const corpId = userInfo?.corpId || 'wwd483c2fba24ae30a'
// 动态获取最后一条学员消息 // 动态获取最后一条学员消息(直接从 archive_messages 表查询,延迟更低)
const fetchLastCustomerMessage = async () => { const fetchLastCustomerMessage = async () => {
if (!customerId || customerId === 'wx_001') return if (!customerId || customerId === 'wx_001') return
try { try {
const res = await fetch(`/api/v1/conversations?customerId=${customerId}&corpId=${corpId}`) const res = await fetch(`/api/v1/conversations/last-message?customerId=${customerId}&corpId=${corpId}`)
const data = await res.json() const data = await res.json()
if (data.code === 0 && data.data && data.data.length > 0) { if (data.code === 0 && data.data && data.data.content) {
// 取最新的会话 setCustomerMsg(data.data.content)
const sessions = data.data as Array<{ sessionId: string; startTime: string }>
const latestSession = sessions.sort(
(a, b) => new Date(b.startTime).getTime() - new Date(a.startTime).getTime()
)[0]
const turnsRes = await fetch(`/api/v1/conversations/${latestSession.sessionId}/turns`)
const turnsData = await turnsRes.json()
if (turnsData.code === 0 && turnsData.data && turnsData.data.length > 0) {
const turns = turnsData.data as ConversationTurn[]
// 取最后一条有内容的学员消息
for (let i = turns.length - 1; i >= 0; i--) {
if (turns[i].studentContent?.trim()) {
setCustomerMsg(turns[i].studentContent.trim())
break
}
}
}
} }
} catch (e) { } catch (e) {
console.error('获取最后一条学员消息失败:', e) console.error('获取最后一条学员消息失败:', e)