979 lines
30 KiB
Go
979 lines
30 KiB
Go
package biz
|
||
|
||
import (
|
||
"ai_scheduler/internal/config"
|
||
dbmodel "ai_scheduler/internal/data/model"
|
||
"ai_scheduler/internal/data/mongo_model"
|
||
"ai_scheduler/internal/entitys"
|
||
"ai_scheduler/internal/pkg"
|
||
"ai_scheduler/internal/pkg/wx"
|
||
"bytes"
|
||
"context"
|
||
"encoding/json"
|
||
"errors"
|
||
"fmt"
|
||
"math/rand"
|
||
"strings"
|
||
"time"
|
||
|
||
"go.mongodb.org/mongo-driver/bson"
|
||
"go.mongodb.org/mongo-driver/mongo/options"
|
||
|
||
"github.com/gofiber/fiber/v2/log"
|
||
"github.com/google/uuid"
|
||
)
|
||
|
||
// WsBroadcaster WebSocket 广播接口(由 services/advice.WxWsHub 实现)
|
||
type WsBroadcaster interface {
|
||
Broadcast(appId string, event string, data interface{})
|
||
}
|
||
|
||
// AdviceWxBiz 个人微信消息回调接入与消息流水
|
||
// 上游(WTAPI)消息通过 login/setCallback 配置的回调地址推送到本服务,
|
||
// 本模块负责容错解析回调报文、记录消息流水、更新客户互动时间。
|
||
type AdviceWxBiz struct {
|
||
cfg *config.Config
|
||
mongo *pkg.Mongo
|
||
wxMsgMongo *mongo_model.AdvicerWxMsgMongo
|
||
clientBiz *AdviceClientBiz
|
||
strategyBiz *AdviceStrategyBiz
|
||
wsBroadcaster WsBroadcaster
|
||
adviceCustomerBiz *AdviceCustomerBiz
|
||
adviceLabelBiz *AdviceLabelBiz
|
||
adviceAdvicerBiz *AdviceAdvicerBiz
|
||
adviceChatBiz *AdviceChatBiz
|
||
adviceProjectBiz *AdviceProjectBiz
|
||
adviceSkillBiz *AdviceSkillBiz
|
||
wxSendBiz *AdviceWxSendBiz
|
||
}
|
||
|
||
func NewAdviceWxBiz(
|
||
cfg *config.Config,
|
||
mongo *pkg.Mongo,
|
||
wxMsgMongo *mongo_model.AdvicerWxMsgMongo,
|
||
clientBiz *AdviceClientBiz,
|
||
strategyBiz *AdviceStrategyBiz,
|
||
wsBroadcaster WsBroadcaster,
|
||
adviceCustomerBiz *AdviceCustomerBiz,
|
||
adviceLabelBiz *AdviceLabelBiz,
|
||
adviceAdvicerBiz *AdviceAdvicerBiz,
|
||
adviceChatBiz *AdviceChatBiz,
|
||
adviceProjectBiz *AdviceProjectBiz,
|
||
adviceSkillBiz *AdviceSkillBiz,
|
||
wxSendBiz *AdviceWxSendBiz,
|
||
) *AdviceWxBiz {
|
||
return &AdviceWxBiz{
|
||
cfg: cfg,
|
||
mongo: mongo,
|
||
wxMsgMongo: wxMsgMongo,
|
||
clientBiz: clientBiz,
|
||
strategyBiz: strategyBiz,
|
||
wsBroadcaster: wsBroadcaster,
|
||
adviceCustomerBiz: adviceCustomerBiz,
|
||
adviceLabelBiz: adviceLabelBiz,
|
||
adviceAdvicerBiz: adviceAdvicerBiz,
|
||
adviceChatBiz: adviceChatBiz,
|
||
adviceProjectBiz: adviceProjectBiz,
|
||
adviceSkillBiz: adviceSkillBiz,
|
||
wxSendBiz: wxSendBiz,
|
||
}
|
||
}
|
||
|
||
// ==================== 回调处理 ====================
|
||
|
||
// HandleCallback 处理上游推送的回调报文(支持单条/数组/字符串包裹 JSON)
|
||
func (a *AdviceWxBiz) HandleCallback(ctx context.Context, body []byte) (err error) {
|
||
body = bytes.TrimSpace(body)
|
||
if len(body) == 0 {
|
||
return nil
|
||
}
|
||
// 数组报文:逐条处理,单条失败不影响其他
|
||
if body[0] == '[' {
|
||
var arr []json.RawMessage
|
||
if e := json.Unmarshal(body, &arr); e == nil {
|
||
for _, item := range arr {
|
||
if e := a.handleOneEvent(ctx, item); e != nil {
|
||
err = e
|
||
}
|
||
}
|
||
return err
|
||
}
|
||
}
|
||
return a.handleOneEvent(ctx, body)
|
||
}
|
||
|
||
// handleOneEvent 处理单条回调事件(基于 callback.md 文档的类型化解析)
|
||
func (a *AdviceWxBiz) handleOneEvent(ctx context.Context, body []byte) error {
|
||
event, parseErr := wx.ParseCallbackEvent(body)
|
||
if parseErr != nil {
|
||
return fmt.Errorf("解析回调报文失败: %w", parseErr)
|
||
}
|
||
|
||
// 按 TypeName 分发
|
||
isAddMsg, isModContacts, isDelContacts, isOffline, isFinder := wx.DispatchCallbackEvent(event)
|
||
|
||
// 联系人变动/离线/视频号事件不处理(仅处理 AddMsg)
|
||
if !isAddMsg {
|
||
_ = isModContacts
|
||
_ = isDelContacts
|
||
_ = isOffline
|
||
_ = isFinder
|
||
return nil
|
||
}
|
||
|
||
if event.Data == nil {
|
||
return nil
|
||
}
|
||
|
||
// 判定消息方向与对话方 wxid
|
||
direction := mongo_model.WxMsgDirectionCustomer
|
||
peerWxid := event.PeerWxid()
|
||
to := event.Data.ToUserName.Val()
|
||
|
||
// 销售自己发的消息通常标记为 Self,但发给 filehelper 的除外:
|
||
// 给文件传输助手发消息时,保持 Customer 方向,让托管流程能作为客户消息触发 AI 回复
|
||
if event.IsSelfMsg() && to != "filehelper" {
|
||
direction = mongo_model.WxMsgDirectionSelf
|
||
} else if len(event.Wxid) == 0 && a.clientBiz != nil {
|
||
// 报文未提供登录微信标识:借助客户绑定库判定
|
||
from := event.Data.FromUserName.Val()
|
||
if _, found, _ := a.clientBiz.FindByWxid(ctx, from); !found {
|
||
if _, foundTo, _ := a.clientBiz.FindByWxid(ctx, to); foundTo {
|
||
direction = mongo_model.WxMsgDirectionSelf
|
||
peerWxid = to
|
||
}
|
||
}
|
||
}
|
||
if len(peerWxid) == 0 {
|
||
peerWxid = event.Data.FromUserName.Val()
|
||
}
|
||
|
||
// 群消息:提取真实发送者
|
||
content := event.Data.Content.Val()
|
||
var senderWxid string
|
||
if event.IsGroupMsg() && direction == mongo_model.WxMsgDirectionCustomer {
|
||
realSender := event.GroupRealSender()
|
||
if realSender != "" {
|
||
senderWxid = realSender
|
||
// 群消息 Content 格式为 "wxid_xxx:\n消息内容",去掉前缀取正文
|
||
if idx := strings.Index(content, ":\n"); idx > 0 {
|
||
content = content[idx+2:]
|
||
}
|
||
}
|
||
}
|
||
|
||
msgType := event.Data.MsgType
|
||
newMsgId := event.Data.NewMsgId.String()
|
||
msgId := event.Data.MsgId.String()
|
||
|
||
// 过滤:公众号消息(gh_ 开头)不存储不展示
|
||
if strings.HasPrefix(peerWxid, "gh_") {
|
||
return nil
|
||
}
|
||
// 过滤:msgType 为 other 的不存储不展示
|
||
if wx.MsgTypeCategory(msgType) == "other" {
|
||
return nil
|
||
}
|
||
|
||
// 媒体消息存原始 XML(去前缀后的),供前端调用下载 API
|
||
var mediaXml string
|
||
if wx.IsMediaMsg(msgType) {
|
||
mediaXml = content
|
||
}
|
||
|
||
record := &mongo_model.AdvicerWxMsgMongo{
|
||
AppId: event.AppId,
|
||
Wxid: peerWxid,
|
||
SelfWxid: event.Wxid,
|
||
SenderWxid: senderWxid,
|
||
Direction: direction,
|
||
MsgType: wx.MsgTypeCategory(msgType),
|
||
Content: content,
|
||
MediaXml: mediaXml,
|
||
MsgId: msgId,
|
||
NewMsgId: newMsgId,
|
||
Source: mongo_model.WxMsgSourceCallback,
|
||
Raw: string(body),
|
||
CreateAt: callbackCreateTime(event.Data.CreateTime),
|
||
}
|
||
|
||
// 消息去重:使用 AppId + NewMsgId(文档推荐去重键)
|
||
dedupKey := event.DedupKey()
|
||
if len(dedupKey) > 1 {
|
||
count, e := a.mongo.Co(a.wxMsgMongo).CountDocuments(ctx, bson.M{"newMsgId": newMsgId, "appId": event.AppId})
|
||
if e != nil {
|
||
return fmt.Errorf("消息去重查询失败: %w", e)
|
||
}
|
||
if count > 0 {
|
||
return nil
|
||
}
|
||
} else if len(msgId) != 0 {
|
||
// 回退:旧格式无 NewMsgId,使用 MsgId 去重
|
||
count, e := a.mongo.Co(a.wxMsgMongo).CountDocuments(ctx, bson.M{"msgId": msgId})
|
||
if e != nil {
|
||
return fmt.Errorf("消息去重查询失败: %w", e)
|
||
}
|
||
if count > 0 {
|
||
return nil
|
||
}
|
||
}
|
||
|
||
if _, e := a.mongo.Co(a.wxMsgMongo).InsertOne(ctx, record); e != nil {
|
||
return fmt.Errorf("消息流水落库失败: %w", e)
|
||
}
|
||
|
||
// WebSocket 实时广播:向对应 appId 的已连接前端推送新消息事件
|
||
if a.wsBroadcaster != nil {
|
||
a.wsBroadcaster.Broadcast(event.AppId, "new_message", map[string]interface{}{
|
||
"appId": event.AppId,
|
||
"wxid": peerWxid,
|
||
"senderWxid": senderWxid,
|
||
"direction": string(direction),
|
||
"msgType": wx.MsgTypeCategory(msgType),
|
||
"content": content,
|
||
"mediaXml": mediaXml,
|
||
"msgId": msgId,
|
||
"createAt": record.CreateAt.Unix(),
|
||
})
|
||
}
|
||
|
||
// 已绑定客户则同步互动时间与消息数
|
||
if a.clientBiz != nil && len(peerWxid) != 0 {
|
||
_ = a.clientBiz.TouchInteraction(ctx, peerWxid, direction)
|
||
}
|
||
|
||
// 智能策略:客户文本消息异步触发 AI 回复决策
|
||
|
||
// 托管自动注册 + AI 回复:双向触发(不限 direction),文本私聊消息检查
|
||
if wx.IsTextMsg(msgType) && !event.IsGroupMsg() {
|
||
appId, selfWxid, msgContent := event.AppId, event.Wxid, content
|
||
|
||
go a.tryHostingAutoRegis(context.Background(), appId, selfWxid, peerWxid, to, msgContent, direction)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// ==================== 托管自动注册 ====================
|
||
|
||
// tryHostingAutoRegis 托管自动回复流程:
|
||
// 1. 查 advicer,确认处于托管状态(hosting_enabled=1)
|
||
// 2. 查 projectInfo 取项目级 wxToken
|
||
// 3. 若回调数据 ToUserName.string == "filehelper" 且 advicer.reply_filehelper=1 → 跳过标签检查,直接进入回复流程
|
||
// 否则调微信 contacts/getDetailInfo 取好友的 labelList("1,2,3" 形式)
|
||
// labelList 为空 → 直接 return,不做回复
|
||
// 4. 取 ai_advice_label 中该销售 hosting=1 的标签 ID 列表
|
||
// 5. 判断好友 labelList 与托管标签是否有交集,无交集 → return
|
||
// 6. 查 ai_advice_customer 取 session_id(无则自动注册;filehelper 不在 customer 表,每次都自动注册)
|
||
// 7. 调 Chat 生成 AI 回复 → hostingChatAndSend 发送
|
||
func (a *AdviceWxBiz) tryHostingAutoRegis(ctx context.Context, appId, selfWxid, peerWxid, toUserName, content string, direction string) {
|
||
if appId == "" || selfWxid == "" || peerWxid == "" {
|
||
return
|
||
}
|
||
|
||
// 1. 通过 appId(即 wx_device_id)查 advicer,确认托管状态
|
||
advicer, err := a.adviceAdvicerBiz.FindByWxDeviceId(ctx, appId)
|
||
if err != nil {
|
||
return
|
||
}
|
||
if advicer.AdvicerID == 0 || advicer.HostingEnabled != 1 {
|
||
return
|
||
}
|
||
|
||
// 2. 查项目信息(含项目级 wxToken),后续 getDetailInfo 与 Chat 都要用
|
||
projectInfo, err := a.adviceProjectBiz.Info(ctx, &entitys.AdvicerProjectInfoReq{ProjectId: advicer.ProjectID})
|
||
if err != nil {
|
||
return
|
||
}
|
||
wxToken := projectInfo.Base.WxToken
|
||
if wxToken == "" {
|
||
return
|
||
}
|
||
|
||
// 3. 文件传输助手特殊路径:仅看回调 ToUserName.string 是否为 "filehelper"
|
||
// 是 → 跳过标签检查,直接检查 advicer.reply_filehelper 开关
|
||
// 否 → 走正常标签交集检查流程
|
||
isFilehelper := toUserName == "filehelper"
|
||
|
||
var sessionId string
|
||
var cachedSession *dbmodel.AiAdviceSession
|
||
var cachedModelInfo *dbmodel.AiAdviceModelSup
|
||
var customerNickName string
|
||
if !isFilehelper {
|
||
// 3a. 调微信 contacts/getDetailInfo 取好友的 labelList
|
||
var detailRes []wx.GetDetailInfoResData
|
||
if err := wx.Request(ctx, wxToken, wx.GetDetailInfo, wx.GetDetailInfoReq{
|
||
AppId: appId,
|
||
Wxids: []string{peerWxid},
|
||
}, &detailRes); err != nil {
|
||
return
|
||
}
|
||
if len(detailRes) == 0 {
|
||
return
|
||
}
|
||
detail := detailRes[0]
|
||
friendLabelList := strings.TrimSpace(detail.LabelList)
|
||
if friendLabelList == "" {
|
||
// 好友没有任何标签,不进入托管回复
|
||
return
|
||
}
|
||
|
||
// 4. 取该销售下 hosting=1 的托管标签 ID 列表
|
||
hostingLabelIds, _ := a.adviceLabelBiz.GetHostingLabelIds(ctx, selfWxid)
|
||
if len(hostingLabelIds) == 0 {
|
||
return
|
||
}
|
||
|
||
// 5. 判断好友标签与托管标签是否有交集
|
||
friendLabelSet := make(map[int]bool)
|
||
for _, s := range strings.Split(friendLabelList, ",") {
|
||
var id int
|
||
if _, e := fmt.Sscanf(strings.TrimSpace(s), "%d", &id); e == nil && id > 0 {
|
||
friendLabelSet[id] = true
|
||
}
|
||
}
|
||
hasIntersection := false
|
||
for _, id := range hostingLabelIds {
|
||
if friendLabelSet[id] {
|
||
hasIntersection = true
|
||
break
|
||
}
|
||
}
|
||
if !hasIntersection {
|
||
return
|
||
}
|
||
|
||
// 6. 查 ai_advice_customer 取 session_id
|
||
customer, err := a.adviceCustomerBiz.FindBySelfWxidAndUserName(ctx, selfWxid, peerWxid)
|
||
if err != nil {
|
||
return
|
||
}
|
||
if customer.CustomerID == 0 {
|
||
return
|
||
}
|
||
|
||
sessionId = customer.SessionId
|
||
// 优先使用微信备注/昵称(更像真人),回退到 customer 表中的昵称
|
||
customerNickName = detail.Remark
|
||
if customerNickName == "" {
|
||
customerNickName = detail.NickName
|
||
}
|
||
if customerNickName == "" {
|
||
customerNickName = customer.NickName
|
||
}
|
||
|
||
// 无会话 → 自动注册
|
||
if sessionId == "" {
|
||
sessionId, cachedSession, cachedModelInfo, err = a.autoRegisSession(ctx, advicer, projectInfo, customer.CustomerID, customerNickName, peerWxid)
|
||
if err != nil {
|
||
// autoRegisSession error
|
||
return
|
||
}
|
||
} else {
|
||
// 复用已有会话时,刷新 mission 确保人格和项目信息是最新的
|
||
}
|
||
} else {
|
||
// filehelper 路径:跳过标签检查,直接检查 advicer.reply_filehelper 开关
|
||
if advicer.ReplyFilehelper != 1 {
|
||
// 开关未开启,不回复 filehelper
|
||
return
|
||
}
|
||
// filehelper 与普通客户同等对待:查 ai_advice_customer 表,复用 session_id
|
||
customerNickName = "文件传输助手"
|
||
customer, err := a.adviceCustomerBiz.FindBySelfWxidAndUserName(ctx, selfWxid, peerWxid)
|
||
if err != nil {
|
||
return
|
||
}
|
||
if customer.CustomerID == 0 {
|
||
if err = a.adviceCustomerBiz.Add(ctx, &entitys.AdvicerCustomerAddReq{
|
||
SelfWxid: selfWxid,
|
||
UserName: peerWxid,
|
||
NickName: customerNickName,
|
||
}); err != nil {
|
||
return
|
||
}
|
||
customer, err = a.adviceCustomerBiz.FindBySelfWxidAndUserName(ctx, selfWxid, peerWxid)
|
||
if err != nil || customer.CustomerID == 0 {
|
||
return
|
||
}
|
||
}
|
||
sessionId = customer.SessionId
|
||
if sessionId == "" {
|
||
sessionId, cachedSession, cachedModelInfo, err = a.autoRegisSession(ctx, advicer, projectInfo, customer.CustomerID, customerNickName, peerWxid)
|
||
if err != nil {
|
||
// filehelper autoRegis error
|
||
return
|
||
}
|
||
} else {
|
||
// 复用已有会话时,刷新 mission 确保人格和项目信息是最新的
|
||
}
|
||
}
|
||
|
||
// 7. 有会话 + 客户发的消息 → 调 Chat 生成 AI 回复并发送
|
||
if sessionId != "" && direction == mongo_model.WxMsgDirectionCustomer && content != "" {
|
||
a.hostingChatAndSend(ctx, sessionId, content, appId, peerWxid, wxToken, cachedSession, cachedModelInfo)
|
||
}
|
||
}
|
||
|
||
// 确保旧 session 也能拿到最新的人格和项目数据,而不是沿用老的空 mission。
|
||
|
||
// autoRegisSession 自动注册托管会话,返回 sessionId + session对象 + modelInfo(供后续 Chat 复用)。
|
||
// advicer 和 projectInfo 由调用方传入,避免重复查询。
|
||
func (a *AdviceWxBiz) autoRegisSession(ctx context.Context,
|
||
advicer dbmodel.AiAdviceAdvicer,
|
||
projectInfo *entitys.AdvicerProjectInfoRes,
|
||
customerId int32, customerNickName string,
|
||
peerWxid string,
|
||
) (string, *dbmodel.AiAdviceSession, *dbmodel.AiAdviceModelSup, error) {
|
||
// 解析版本
|
||
versionId := advicer.HostingVersionId
|
||
var err error
|
||
if versionId == "" {
|
||
versionId, _, err = a.adviceAdvicerBiz.GetLatestVersion(ctx, advicer.AdvicerID)
|
||
if err != nil {
|
||
return "", nil, nil, fmt.Errorf("取最新版本版本失败: %w", err)
|
||
}
|
||
}
|
||
|
||
if len(projectInfo.ModelInfo.Key) == 0 || len(projectInfo.ModelInfo.ChatModel) == 0 {
|
||
return "", nil, nil, fmt.Errorf("项目未配置模型 projectId=%d", advicer.ProjectID)
|
||
}
|
||
|
||
// 查模型信息(用于后续 Chat,通过 projectInfo 中的 ModelSupID 获取)
|
||
modelInfo, err := a.adviceProjectBiz.ModelInfo(projectInfo.Base.ModelSupID)
|
||
if err != nil {
|
||
return "", nil, nil, fmt.Errorf("查模型信息失败: %w", err)
|
||
}
|
||
|
||
// 查版本数据
|
||
versionInfo, err := a.adviceAdvicerBiz.VersionInfo(ctx, &entitys.AdvicerVersionInfoReq{Id: versionId})
|
||
if err != nil {
|
||
return "", nil, nil, fmt.Errorf("查版本数据失败: %w", err)
|
||
}
|
||
|
||
// 加载销售技巧
|
||
var talkSkillEntity *mongo_model.AdvicerTalkSkillMongoEntity
|
||
if advicer.HostingSkillId != "" {
|
||
// 指定了技巧ID,直接加载
|
||
if skillInfo, e := a.adviceSkillBiz.Info(ctx, &entitys.AdvicerTalkSkillInfoReq{Id: advicer.HostingSkillId}); e == nil {
|
||
talkSkillEntity = skillInfo.Entity()
|
||
}
|
||
} else {
|
||
// 未指定技巧ID,取该项目最新的技巧
|
||
if list, e := a.adviceSkillBiz.VersionList(ctx, &entitys.AdvicerTalkSkillListReq{ProjectId: projectInfo.Base.ProjectID}); e == nil && len(list) > 0 {
|
||
talkSkillEntity = list[0].AdvicerTalkSkillMongo.Entity()
|
||
}
|
||
}
|
||
// 加载客户信息(MongoDB advicer_client)
|
||
var clientEntity *mongo_model.AdvicerClientMongoEntity
|
||
if peerWxid != "" {
|
||
var clientDoc mongo_model.AdvicerClientMongo
|
||
filter := bson.M{"projectId": projectInfo.Base.ProjectID, "advicerId": advicer.AdvicerID, "wxid": peerWxid}
|
||
if e := a.mongo.Co(mongo_model.NewAdvicerClientMongo()).FindOne(ctx, filter).Decode(&clientDoc); e == nil {
|
||
clientEntity = clientDoc.Entity()
|
||
}
|
||
}
|
||
// 构建完整 ProjectInfo(确保 MySQL 项目名不丢)
|
||
projEntity := projectInfo.ConfigInfo.Entity()
|
||
if projEntity.ProjectInfo.Name == "" && projectInfo.Base.Name != "" {
|
||
projEntity.ProjectInfo.Name = projectInfo.Base.Name
|
||
}
|
||
|
||
// 加载项目资料(advicer_project_data 集合,扁平栏目数据)
|
||
var projectData map[string]interface{}
|
||
if projData, e := a.adviceProjectBiz.ProjectDataLoad(ctx, projectInfo.Base.ProjectID); e == nil && len(projData.Data) > 0 {
|
||
projectData = projData.Data
|
||
}
|
||
|
||
// 构建 ChatData(全量数据注入)
|
||
chatData := &entitys.ChatData{
|
||
AdvicerInfo: advicer.Entity(),
|
||
AdvicerVersion: versionInfo.Data,
|
||
ProjectInfo: projEntity,
|
||
ProjectData: projectData,
|
||
RuleDimension: projectInfo.Base.RuleDimension,
|
||
TalkSkill: talkSkillEntity,
|
||
ClientInfo: clientEntity,
|
||
}
|
||
|
||
// 构建注册请求(托管场景额外强调短消息格式)
|
||
mission := fmt.Sprintf("与%s的自动托管会话。注意:每条回复必须简短(10~30字),像真人微信聊天一样拆成多条短消息,用\\n分隔。", customerNickName)
|
||
if customerNickName == "" {
|
||
mission = "自动托管会话。注意:每条回复必须简短(10~30字),像真人微信聊天一样拆成多条短消息,用\\n分隔。"
|
||
}
|
||
// 将项目基本信息拼接到 mission,让后续每次 Chat 的 taskPrompt 都能带上项目名称和地址。
|
||
// 原因:Chat 时只传 taskPrompt + 聊天历史 + 用户消息,PreviousResponseID 续写不会让模型可靠获取项目信息,
|
||
// 因此必须把项目信息直接写入 mission 字段,随 session 持久化,每次 Chat 都能看到。
|
||
// 项目信息:优先用 MySQL 的 Base.Name(一定有),MongoDB 的 ProjectInfo 作为补充
|
||
projectDesc := "\n[项目信息]"
|
||
hasProjectInfo := false
|
||
if projectInfo.Base.Name != "" {
|
||
projectDesc += "\n项目名称:" + projectInfo.Base.Name
|
||
hasProjectInfo = true
|
||
}
|
||
if pi := projectInfo.ConfigInfo.ProjectInfo; pi.Address != "" {
|
||
projectDesc += "\n项目地址:" + pi.Address
|
||
hasProjectInfo = true
|
||
}
|
||
if hasProjectInfo {
|
||
mission = mission + projectDesc
|
||
}
|
||
// 将销售人格拼接到 mission
|
||
if persona := formatPersonaPrompt(chatData.AdvicerVersion); persona != "" {
|
||
mission = mission + "\n\n" + persona
|
||
}
|
||
regisReq := &entitys.AdvicerChatRegistReq{
|
||
AdvicerVersionId: versionId,
|
||
TalkSkillId: advicer.HostingSkillId,
|
||
Mission: mission,
|
||
AdvicerId: advicer.AdvicerID,
|
||
}
|
||
|
||
// 构造 session 对象(Regis 内部会写入相同数据,这里同步保存供后续 Chat 复用)
|
||
sessionId := uuid.New().String()
|
||
session := &dbmodel.AiAdviceSession{
|
||
SessionID: sessionId,
|
||
ProjectID: projectInfo.Base.ProjectID,
|
||
SupID: projectInfo.Base.ModelSupID,
|
||
AdvicerVersionID: versionId,
|
||
AdvicerID: advicer.AdvicerID,
|
||
TalkSkillID: advicer.HostingSkillId,
|
||
Mission: mission,
|
||
}
|
||
|
||
// 注册会话
|
||
regisSessionId, err := a.adviceChatBiz.Regis(ctx, chatData, regisReq, projectInfo)
|
||
if err != nil {
|
||
return "", nil, nil, fmt.Errorf("注册会话失败: %w", err)
|
||
}
|
||
sessionId = regisSessionId
|
||
session.SessionID = sessionId
|
||
|
||
// 回写 customer.session_id
|
||
if err = a.adviceCustomerBiz.UpdateSessionId(ctx, customerId, sessionId); err != nil {
|
||
// ignore
|
||
}
|
||
return sessionId, session, &modelInfo, nil
|
||
}
|
||
|
||
// hostingChatAndSend 调用 Chat 生成 AI 回复,并通过微信发送。
|
||
// 当 session 和 modelInfo 非空时(新注册路径),使用 ChatWithSession 跳过重复查询。
|
||
func (a *AdviceWxBiz) hostingChatAndSend(ctx context.Context, sessionId, content, appId, peerWxid, wxToken string,
|
||
session *dbmodel.AiAdviceSession, modelInfo *dbmodel.AiAdviceModelSup,
|
||
) {
|
||
var assistant mongo_model.Assistant
|
||
var err error
|
||
|
||
chatReq := &entitys.AdvicerChatReq{SessionId: sessionId, Content: content}
|
||
if session != nil && session.SessionID != "" && modelInfo != nil && modelInfo.SupID != 0 {
|
||
// 新注册路径:session/model 已在 autoRegisSession 中获取,跳过重复查询
|
||
assistant, err = a.adviceChatBiz.ChatWithSession(ctx, chatReq, session, modelInfo)
|
||
} else {
|
||
// 已有会话路径:走原 Chat 查询
|
||
assistant, err = a.adviceChatBiz.Chat(ctx, chatReq)
|
||
}
|
||
if err != nil {
|
||
return
|
||
}
|
||
reply := strings.TrimSpace(assistant.Result)
|
||
if reply == "" {
|
||
return
|
||
}
|
||
// 拆分回复为多段消息,模拟真人分段发送
|
||
chunks := splitReplyToChunks(reply)
|
||
// 模拟阅读客户消息 + 思考的延迟(3~8秒)
|
||
readDelay := time.Duration(3000+rand.Intn(5001)) * time.Millisecond
|
||
time.Sleep(readDelay)
|
||
// 通过微信分段发送 AI 回复(使用项目级 wx_token)
|
||
if a.wxSendBiz != nil {
|
||
if _, err = a.wxSendBiz.SendMultiWithToken(ctx, wxToken, appId, peerWxid, chunks, mongo_model.WxMsgSourceCallback); err != nil {
|
||
// ignore
|
||
} else {
|
||
log.Infof("[微信回复] sessionId=%s reply=%s", sessionId, reply)
|
||
}
|
||
}
|
||
}
|
||
|
||
// splitReplyToChunks 将 AI 回复拆分为多条短消息,模拟真人分段发送。
|
||
// 拆分策略:先按换行分段,过长段落按句末标点拆分,过短段落合并。
|
||
func splitReplyToChunks(reply string) []string {
|
||
const maxChunkLen = 80 // 每条消息最大字数(接近真人微信习惯)
|
||
const minChunkLen = 15 // 低于此字数考虑合并到上一条
|
||
|
||
// 1. 按双换行或单换行拆分为段落
|
||
rawParagraphs := strings.Split(reply, "\n\n")
|
||
var paragraphs []string
|
||
for _, p := range rawParagraphs {
|
||
p = strings.TrimSpace(p)
|
||
if p == "" {
|
||
continue
|
||
}
|
||
// 段落内如果还有单换行,也拆开
|
||
for _, line := range strings.Split(p, "\n") {
|
||
line = strings.TrimSpace(line)
|
||
if line != "" {
|
||
paragraphs = append(paragraphs, line)
|
||
}
|
||
}
|
||
}
|
||
if len(paragraphs) == 0 {
|
||
return []string{reply}
|
||
}
|
||
|
||
// 2. 过长段落按句末标点拆分
|
||
sentenceEnds := []rune{'。', '!', '?', ';', '.', '!', '?', ';'}
|
||
var segments []string
|
||
for _, p := range paragraphs {
|
||
if len([]rune(p)) <= maxChunkLen {
|
||
segments = append(segments, p)
|
||
continue
|
||
}
|
||
// 按句末标点拆分
|
||
runes := []rune(p)
|
||
start := 0
|
||
for i, r := range runes {
|
||
for _, sep := range sentenceEnds {
|
||
if r == sep && i-start >= 10 {
|
||
seg := string(runes[start : i+1])
|
||
if len([]rune(seg)) <= maxChunkLen {
|
||
segments = append(segments, seg)
|
||
start = i + 1
|
||
}
|
||
break
|
||
}
|
||
}
|
||
}
|
||
// 剩余部分
|
||
if start < len(runes) {
|
||
remaining := strings.TrimSpace(string(runes[start:]))
|
||
if remaining != "" {
|
||
segments = append(segments, remaining)
|
||
}
|
||
}
|
||
}
|
||
|
||
// 3. 合并过短的段落(模拟真人一口气发完短句)
|
||
var chunks []string
|
||
buf := ""
|
||
for _, s := range segments {
|
||
if buf == "" {
|
||
buf = s
|
||
} else if len([]rune(buf))+len([]rune(s))+1 <= maxChunkLen && len([]rune(s)) < minChunkLen {
|
||
// 当前段很短,合并到上一条
|
||
buf = buf + "\n" + s
|
||
} else {
|
||
chunks = append(chunks, buf)
|
||
buf = s
|
||
}
|
||
}
|
||
if buf != "" {
|
||
chunks = append(chunks, buf)
|
||
}
|
||
|
||
if len(chunks) == 0 {
|
||
return []string{reply}
|
||
}
|
||
return chunks
|
||
}
|
||
|
||
// callbackCreateTime 将回调时间戳转为 time.Time(兼容秒/毫秒,缺省取当前时间)
|
||
func callbackCreateTime(ts int64) time.Time {
|
||
if ts <= 0 {
|
||
return time.Now()
|
||
}
|
||
if ts > 1e12 { // 毫秒级
|
||
return time.UnixMilli(ts)
|
||
}
|
||
return time.Unix(ts, 0)
|
||
}
|
||
|
||
// VerifyCallbackToken 校验回调请求合法性。
|
||
// 上游回调若携带 token(query 或 body 字段)则与配置比对;未携带时放行。
|
||
func (a *AdviceWxBiz) VerifyCallbackToken(query map[string]string, body []byte) bool {
|
||
expect := a.cfg.Advicer.CallbackToken
|
||
if len(expect) == 0 {
|
||
return true
|
||
}
|
||
got := query["token"]
|
||
if got == "" {
|
||
// 从 body 中提取 token 字段(兼容多种命名)
|
||
m := toMap(body)
|
||
if m != nil {
|
||
for _, k := range []string{"token", "Token", "verifyToken", "verify_token"} {
|
||
if v, ok := m[k]; ok {
|
||
if s, ok := v.(string); ok && s != "" {
|
||
got = s
|
||
break
|
||
}
|
||
}
|
||
}
|
||
}
|
||
}
|
||
if got == "" {
|
||
return true
|
||
}
|
||
return got == expect
|
||
}
|
||
|
||
// SetCallback 调用上游接口配置消息回调地址(login/setCallback)
|
||
func (a *AdviceWxBiz) SetCallback(ctx context.Context, req *entitys.AdvicerWxCallbackReq) (err error) {
|
||
token := a.cfg.Advicer.WxToken
|
||
if len(token) == 0 {
|
||
return errors.New("未配置 advicer.wx_token,无法调用上游接口")
|
||
}
|
||
if len(req.CallbackUrl) == 0 {
|
||
return errors.New("回调地址不能为空")
|
||
}
|
||
cbToken := a.cfg.Advicer.CallbackToken
|
||
if len(cbToken) == 0 {
|
||
cbToken = token
|
||
}
|
||
var res wx.SetCallbackResData
|
||
if err = wx.Request(ctx, token, wx.SetCallback, wx.SetCallbackReq{
|
||
Token: cbToken,
|
||
CallbackUrl: req.CallbackUrl,
|
||
}, &res); err != nil {
|
||
return fmt.Errorf("设置回调失败: %w", err)
|
||
}
|
||
return nil
|
||
}
|
||
|
||
// ==================== 消息流水查询 ====================
|
||
|
||
// MsgList 查询微信消息流水(按时间倒序)
|
||
func (a *AdviceWxBiz) MsgList(ctx context.Context, param *entitys.AdvicerWxMsgListReq) (list []mongo_model.AdvicerWxMsgItem, err error) {
|
||
filter := bson.M{}
|
||
if len(param.AppId) != 0 {
|
||
filter["appId"] = param.AppId
|
||
}
|
||
if len(param.Wxid) != 0 {
|
||
filter["wxid"] = param.Wxid
|
||
}
|
||
if len(param.Direction) != 0 {
|
||
filter["direction"] = param.Direction
|
||
}
|
||
timeCond := bson.M{}
|
||
if len(param.StartAt) != 0 {
|
||
if t, e := parseTimeStr(param.StartAt); e == nil {
|
||
timeCond["$gte"] = t
|
||
}
|
||
}
|
||
if len(param.EndAt) != 0 {
|
||
if t, e := parseTimeStr(param.EndAt); e == nil {
|
||
timeCond["$lte"] = t
|
||
}
|
||
}
|
||
if len(timeCond) != 0 {
|
||
filter["createAt"] = timeCond
|
||
}
|
||
|
||
opts := options.Find().SetSort(bson.D{{Key: "createAt", Value: -1}})
|
||
if param.PageSize > 0 {
|
||
page := param.Page
|
||
if page < 1 {
|
||
page = 1
|
||
}
|
||
size := int64(param.PageSize)
|
||
opts.SetSkip(int64(page-1) * size).SetLimit(size)
|
||
} else {
|
||
// 未分页时限制最大返回条数,避免全量拉取
|
||
opts.SetLimit(200)
|
||
}
|
||
|
||
cursor, err := a.mongo.Co(a.wxMsgMongo).Find(ctx, filter, opts)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
for cursor.Next(ctx) {
|
||
var item mongo_model.AdvicerWxMsgItem
|
||
if err = cursor.Decode(&item); err != nil {
|
||
return nil, err
|
||
}
|
||
list = append(list, item)
|
||
}
|
||
if err = cursor.Err(); err != nil {
|
||
return nil, err
|
||
}
|
||
return list, nil
|
||
}
|
||
|
||
// ==================== 会话列表与消息记录 ====================
|
||
|
||
// ConversationList 查询会话列表(按最新消息时间倒序,聚合每个 wxid 的最新一条消息)
|
||
func (a *AdviceWxBiz) ConversationList(ctx context.Context, param *entitys.AdvicerWxConversationListReq) ([]entitys.AdvicerWxConversationItem, error) {
|
||
filter := bson.M{}
|
||
if len(param.SelfWxid) != 0 {
|
||
filter["selfWxid"] = param.SelfWxid
|
||
}
|
||
page := param.Page
|
||
if page < 1 {
|
||
page = 1
|
||
}
|
||
size := param.PageSize
|
||
if size <= 0 {
|
||
size = 50
|
||
}
|
||
if size > 200 {
|
||
size = 200
|
||
}
|
||
|
||
pipeline := []bson.M{
|
||
{"$match": filter},
|
||
{"$sort": bson.M{"createAt": -1}},
|
||
{"$group": bson.M{
|
||
"_id": "$wxid",
|
||
"lastMsg": bson.M{"$first": "$content"},
|
||
"lastMsgType": bson.M{"$first": "$msgType"},
|
||
"lastMsgTime": bson.M{"$first": "$createAt"},
|
||
"lastDirection": bson.M{"$first": "$direction"},
|
||
"unreadCount": bson.M{"$sum": bson.M{
|
||
"$cond": []interface{}{
|
||
bson.M{"$and": []interface{}{
|
||
bson.M{"$eq": []interface{}{"$direction", mongo_model.WxMsgDirectionCustomer}},
|
||
bson.M{"$eq": []interface{}{"$read", false}},
|
||
}},
|
||
1, 0,
|
||
},
|
||
}},
|
||
}},
|
||
{"$sort": bson.M{"lastMsgTime": -1}},
|
||
{"$skip": int64(page-1) * int64(size)},
|
||
{"$limit": int64(size)},
|
||
}
|
||
|
||
cursor, err := a.mongo.Co(a.wxMsgMongo).Aggregate(ctx, pipeline)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
var list []entitys.AdvicerWxConversationItem
|
||
for cursor.Next(ctx) {
|
||
var row struct {
|
||
Wxid string `bson:"_id"`
|
||
LastMsg string `bson:"lastMsg"`
|
||
LastMsgType string `bson:"lastMsgType"`
|
||
LastMsgTime time.Time `bson:"lastMsgTime"`
|
||
LastDirection string `bson:"lastDirection"`
|
||
UnreadCount int `bson:"unreadCount"`
|
||
}
|
||
if err = cursor.Decode(&row); err != nil {
|
||
return nil, err
|
||
}
|
||
// 截断过长消息
|
||
preview := row.LastMsg
|
||
if len([]rune(preview)) > 40 {
|
||
preview = string([]rune(preview)[:40]) + "…"
|
||
}
|
||
if row.LastMsgType != "text" {
|
||
preview = "[" + row.LastMsgType + "]"
|
||
}
|
||
list = append(list, entitys.AdvicerWxConversationItem{
|
||
Wxid: row.Wxid,
|
||
LastMsg: preview,
|
||
LastMsgType: row.LastMsgType,
|
||
LastMsgTime: row.LastMsgTime.Unix(),
|
||
LastDirection: row.LastDirection,
|
||
UnreadCount: row.UnreadCount,
|
||
IsGroup: strings.HasSuffix(row.Wxid, "@chatroom"),
|
||
})
|
||
}
|
||
if err = cursor.Err(); err != nil {
|
||
return nil, err
|
||
}
|
||
return list, nil
|
||
}
|
||
|
||
// MarkAsRead 标记某个会话中所有客户消息为已读
|
||
func (a *AdviceWxBiz) MarkAsRead(ctx context.Context, param *entitys.AdvicerWxMarkAsReadReq) (int64, error) {
|
||
if len(param.Wxid) == 0 {
|
||
return 0, nil
|
||
}
|
||
filter := bson.M{
|
||
"wxid": param.Wxid,
|
||
"direction": mongo_model.WxMsgDirectionCustomer,
|
||
"read": false,
|
||
}
|
||
if len(param.SelfWxid) != 0 {
|
||
filter["selfWxid"] = param.SelfWxid
|
||
}
|
||
update := bson.M{"$set": bson.M{"read": true}}
|
||
result, err := a.mongo.Co(a.wxMsgMongo).UpdateMany(ctx, filter, update)
|
||
if err != nil {
|
||
return 0, err
|
||
}
|
||
return result.ModifiedCount, nil
|
||
}
|
||
|
||
// ConversationMsgs 查询某个会话的消息记录(按时间正序,用于聊天展示)
|
||
func (a *AdviceWxBiz) ConversationMsgs(ctx context.Context, param *entitys.AdvicerWxConversationMsgsReq) ([]mongo_model.AdvicerWxMsgItem, error) {
|
||
filter := bson.M{}
|
||
if len(param.SelfWxid) != 0 {
|
||
filter["selfWxid"] = param.SelfWxid
|
||
}
|
||
if len(param.Wxid) != 0 {
|
||
filter["wxid"] = param.Wxid
|
||
}
|
||
opts := options.Find().SetSort(bson.D{{Key: "createAt", Value: 1}})
|
||
if param.PageSize > 0 {
|
||
page := param.Page
|
||
if page < 1 {
|
||
page = 1
|
||
}
|
||
size := int64(param.PageSize)
|
||
if size > 500 {
|
||
size = 500
|
||
}
|
||
opts.SetSkip(int64(page-1) * size).SetLimit(size)
|
||
} else {
|
||
opts.SetLimit(200)
|
||
}
|
||
|
||
cursor, err := a.mongo.Co(a.wxMsgMongo).Find(ctx, filter, opts)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
var list []mongo_model.AdvicerWxMsgItem
|
||
for cursor.Next(ctx) {
|
||
var item mongo_model.AdvicerWxMsgItem
|
||
if err = cursor.Decode(&item); err != nil {
|
||
return nil, err
|
||
}
|
||
list = append(list, item)
|
||
}
|
||
if err = cursor.Err(); err != nil {
|
||
return nil, err
|
||
}
|
||
return list, nil
|
||
}
|
||
|
||
// ==================== 辅助工具 ====================
|
||
|
||
// toMap 将 JSON 解析为 map(数字保留原文,避免大 MsgId 精度丢失)
|
||
func toMap(b []byte) map[string]interface{} {
|
||
dec := json.NewDecoder(bytes.NewReader(b))
|
||
dec.UseNumber()
|
||
var m map[string]interface{}
|
||
if err := dec.Decode(&m); err != nil {
|
||
return nil
|
||
}
|
||
return m
|
||
}
|
||
|
||
// parseTimeStr 容错解析时间字符串(支持多种常见格式)
|
||
func parseTimeStr(s string) (time.Time, error) {
|
||
s = strings.TrimSpace(s)
|
||
if s == "" {
|
||
return time.Time{}, errors.New("空时间")
|
||
}
|
||
layouts := []string{
|
||
time.RFC3339,
|
||
"2006-01-02 15:04:05",
|
||
"2006-01-02 15:04",
|
||
"2006-01-02",
|
||
}
|
||
for _, l := range layouts {
|
||
if t, err := time.ParseInLocation(l, s, time.Local); err == nil {
|
||
return t, nil
|
||
}
|
||
}
|
||
return time.Time{}, fmt.Errorf("无法解析时间: %s", s)
|
||
}
|