ai_scheduler/internal/biz/advice_iterate.go

480 lines
16 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package biz
import (
"ai_scheduler/internal/biz/llm_service/third_party"
"ai_scheduler/internal/config"
"ai_scheduler/internal/data/constants"
"ai_scheduler/internal/data/impl"
"ai_scheduler/internal/data/mongo_model"
"ai_scheduler/internal/entitys"
"ai_scheduler/internal/pkg"
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"time"
"github.com/sashabaranov/go-openai"
"go.mongodb.org/mongo-driver/bson"
"go.mongodb.org/mongo-driver/bson/primitive"
"go.mongodb.org/mongo-driver/mongo/options"
)
// AdviceIterateBiz 对话完成自我迭代
// 销售与客户的对话结束(最后一条消息超过 dialog_idle_minutes,默认 30 分钟)后,
// 将本轮聊天记录交给 LLM 再次分析,提取新信息并新建/完善客户画像。
type AdviceIterateBiz struct {
cfg *config.Config
openai *third_party.OpenAi
mongo *pkg.Mongo
clientBiz *AdviceClientBiz
wxMsgMongo *mongo_model.AdvicerWxMsgMongo
projectImpl *impl.AdviceProjectImpl
modelSupImpl *impl.AiAdviceModelSupImpl
customerImpl *impl.AdviceCustomerImpl
advicerBiz *AdviceAdvicerBiz
followBiz *AdviceFollowBiz
}
func NewAdviceIterateBiz(
cfg *config.Config,
openai *third_party.OpenAi,
mongo *pkg.Mongo,
clientBiz *AdviceClientBiz,
wxMsgMongo *mongo_model.AdvicerWxMsgMongo,
projectImpl *impl.AdviceProjectImpl,
modelSupImpl *impl.AiAdviceModelSupImpl,
customerImpl *impl.AdviceCustomerImpl,
advicerBiz *AdviceAdvicerBiz,
followBiz *AdviceFollowBiz,
) *AdviceIterateBiz {
return &AdviceIterateBiz{
cfg: cfg,
openai: openai,
mongo: mongo,
clientBiz: clientBiz,
wxMsgMongo: wxMsgMongo,
projectImpl: projectImpl,
modelSupImpl: modelSupImpl,
customerImpl: customerImpl,
advicerBiz: advicerBiz,
followBiz: followBiz,
}
}
// Run 执行画像迭代:扫描未分析消息 → 按客户分组 → 对话完成判定 → LLM 提取增量 → 合并画像 → 标记已分析。
// Force=true 忽略对话完成(闲置时长)等待,立即迭代(手动触发用)。
func (a *AdviceIterateBiz) Run(ctx context.Context, param *entitys.AdvicerIterateRunReq) (res *entitys.AdvicerIterateRunRes, err error) {
res = &entitys.AdvicerIterateRunRes{Details: []string{}}
if param == nil {
param = &entitys.AdvicerIterateRunReq{}
}
// 1. 扫描未分析消息(时间正序,单轮上限 500 条)
filter := bson.M{"analyzed": false}
if len(param.ClientId) != 0 {
client, e := a.clientBiz.Info(ctx, &entitys.AdvicerClientInfoReq{Id: param.ClientId})
if e != nil {
return res, fmt.Errorf("客户不存在: %w", e)
}
if len(client.Wxid) == 0 {
return res, errors.New("该客户未绑定微信 wxid,无聊天记录可迭代")
}
filter["wxid"] = client.Wxid
}
cursor, e := a.mongo.Co(a.wxMsgMongo).Find(ctx, filter,
options.Find().SetSort(bson.D{{Key: "createAt", Value: 1}}).SetLimit(500))
if e != nil {
return res, e
}
var msgs []mongo_model.AdvicerWxMsgItem
for cursor.Next(ctx) {
var item mongo_model.AdvicerWxMsgItem
if e = cursor.Decode(&item); e != nil {
return res, e
}
msgs = append(msgs, item)
}
if e = cursor.Err(); e != nil {
return res, e
}
res.Scanned = len(msgs)
if len(msgs) == 0 {
return res, nil
}
// 2. 按客户 wxid 分组(保持时间顺序)
group := map[string][]mongo_model.AdvicerWxMsgItem{}
var order []string
for _, m := range msgs {
if len(m.Wxid) == 0 {
continue
}
if _, ok := group[m.Wxid]; !ok {
order = append(order, m.Wxid)
}
group[m.Wxid] = append(group[m.Wxid], m)
}
idle := time.Duration(a.idleMinutes()) * time.Minute
for _, wxid := range order {
items := group[wxid]
// 2.1 客户绑定校验(未绑定客户不标记,待绑定后再迭代)
client, found, e := a.clientBiz.FindByWxid(ctx, wxid)
if e != nil || !found {
continue
}
// 2.2 对话完成判定:最新一条消息(任意方向)距现在超过闲置时长
if !param.Force {
latest := a.latestMsgTime(ctx, wxid)
if latest.IsZero() || time.Since(latest) < idle {
res.Idle++
continue
}
}
// 2.3 LLM 分析并合并画像
summary, updated, e := a.iterateClient(ctx, &client, items)
// 2.4 标记本轮消息已分析(无论画像是否有增量,避免重复分析)
ids := make([]primitive.ObjectID, 0, len(items))
for _, m := range items {
ids = append(ids, m.Id)
}
_, _ = a.mongo.Co(a.wxMsgMongo).UpdateMany(ctx,
bson.M{"_id": bson.M{"$in": ids}},
bson.M{"$set": bson.M{"analyzed": true}})
if e != nil {
res.Details = append(res.Details, fmt.Sprintf("%s:迭代失败 %v", clientName(&client), e))
continue
}
if updated {
res.Updated++
text := clientName(&client) + ":画像已更新"
if strings.TrimSpace(summary) != "" {
text += "(" + strings.TrimSpace(summary) + ")"
}
res.Details = append(res.Details, text)
} else {
res.Details = append(res.Details, clientName(&client)+":本轮无新增画像信息")
}
}
return res, nil
}
// iterateClient 对单客户执行画像迭代:LLM 提取增量 → 合并画像 → 更新 lastIterateAt
func (a *AdviceIterateBiz) iterateClient(ctx context.Context, c *mongo_model.AdvicerClientItem, newMsgs []mongo_model.AdvicerWxMsgItem) (summary string, updated bool, err error) {
_, model := projectModelOf(ctx, a.projectImpl, a.modelSupImpl, c.ProjectId)
if model == nil {
return "", false, errors.New("项目未配置可用模型(请检查 modelSupId)")
}
history := a.recentMsgs(ctx, c.Wxid, 30)
manual, wxRemark := a.fetchRemarks(ctx, c)
content := buildIterateUserContent(c, history, newMsgs, manual, wxRemark)
messages := []openai.ResponseInputMessage{
{Role: openai.ChatMessageRoleSystem, Content: constants.ClientIteratePrompt},
{Role: openai.ChatMessageRoleUser, Content: content},
}
resp, err := a.openai.CreateResponseMessages(ctx, model.Key, model.URL, model.ChatModel, messages, "")
if err != nil {
return "", false, fmt.Errorf("分析聊天记录失败: %w", err)
}
var patch iteratePatch
if err = json.Unmarshal([]byte(extractJsonObject(resp.GetOutputText())), &patch); err != nil {
return "", false, fmt.Errorf("解析迭代结果失败: %w", err)
}
// 合并动态栏目(增量:以现有画像为基底,LLM 输出逐个覆盖非空栏目)
updated, err = a.applyIterateSections(ctx, c, patch.Sections, false)
if err != nil {
return "", false, err
}
return strings.TrimSpace(patch.Summary), updated, nil
}
// applyIterateSections 把 LLM 输出的栏目合并/替换进画像并落库。
// replace=true 为全量替换(「画像更新」按钮),false 为增量合并(定时迭代)。
func (a *AdviceIterateBiz) applyIterateSections(ctx context.Context, c *mongo_model.AdvicerClientItem, sections map[string]interface{}, replace bool) (updated bool, err error) {
merged := map[string]interface{}{}
if !replace {
for k, v := range c.Profile {
merged[k] = v
}
}
for k, v := range sections {
name := strings.TrimSpace(k)
if name == "" || isEmptySection(v) {
continue
}
merged[name] = v
}
now := time.Now()
set := bson.M{"lastIterateAt": now}
if pkg.JsonStringIgonErr(merged) != pkg.JsonStringIgonErr(c.Profile) {
set["profile"] = merged
set["lastUpdateTime"] = now
updated = true
}
_, err = a.mongo.Co(a.clientBiz.AdvicerClientMongo).UpdateOne(ctx, bson.M{"_id": c.Id}, bson.M{"$set": set})
return
}
// isEmptySection 判断栏目值是否为空(空串/空数组/空对象/nil)。
func isEmptySection(v interface{}) bool {
switch x := v.(type) {
case nil:
return true
case string:
return strings.TrimSpace(x) == ""
case []interface{}:
return len(x) == 0
case map[string]interface{}:
return len(x) == 0
}
return false
}
// fetchRemarks 取客户备注:人工备注(description) + 微信备注名(remark),均存于 MySQL ai_advice_customer。
// 靠 appId 反查销售 self_wxid,再按 (self_wxid, 客户wxid) 定位客户记录;查不到返回空串不报错。
func (a *AdviceIterateBiz) fetchRemarks(ctx context.Context, c *mongo_model.AdvicerClientItem) (manual, wxRemark string) {
if a.customerImpl == nil || a.advicerBiz == nil || c == nil || len(c.AppId) == 0 || len(c.Wxid) == 0 {
return "", ""
}
adv, e := a.advicerBiz.FindByWxDeviceId(ctx, c.AppId)
if e != nil || len(adv.WxId) == 0 {
return "", ""
}
cust, e := a.customerImpl.FindBySelfWxidAndUserName(ctx, adv.WxId, c.Wxid)
if e != nil {
return "", ""
}
return strings.TrimSpace(cust.Description), strings.TrimSpace(cust.Remark)
}
// idleMinutes 对话完成闲置时长(分钟,默认 30)
func (a *AdviceIterateBiz) idleMinutes() int {
d := a.cfg.Advicer.DialogIdleMinutes
if d <= 0 {
d = 30
}
return d
}
// latestMsgTime 客户最新一条消息时间(任意方向)
func (a *AdviceIterateBiz) latestMsgTime(ctx context.Context, wxid string) time.Time {
res := a.mongo.Co(a.wxMsgMongo).FindOne(ctx, bson.M{"wxid": wxid},
options.FindOne().SetSort(bson.D{{Key: "createAt", Value: -1}}))
if res.Err() != nil {
return time.Time{}
}
var m mongo_model.AdvicerWxMsgMongo
if err := res.Decode(&m); err != nil {
return time.Time{}
}
return m.CreateAt
}
// recentMsgs 拉取客户最近聊天记录(时间正序返回)
func (a *AdviceIterateBiz) recentMsgs(ctx context.Context, wxid string, limit int64) (list []mongo_model.AdvicerWxMsgMongo) {
if len(wxid) == 0 {
return nil
}
cursor, err := a.mongo.Co(a.wxMsgMongo).Find(ctx, bson.M{"wxid": wxid},
options.Find().SetSort(bson.D{{Key: "createAt", Value: -1}}).SetLimit(limit))
if err != nil {
return nil
}
for cursor.Next(ctx) {
var m mongo_model.AdvicerWxMsgMongo
if err := cursor.Decode(&m); err != nil {
return nil
}
list = append(list, m)
}
for i, j := 0, len(list)-1; i < j; i, j = i+1, j-1 {
list[i], list[j] = list[j], list[i]
}
return list
}
// ==================== 迭代辅助类型与函数 ====================
// iteratePatch 画像迭代/生成结果(sections = 栏目名→内容;summary = 本次说明)
type iteratePatch struct {
Sections map[string]interface{} `json:"sections"`
Summary string `json:"summary"`
}
// RegenerateProfile 画像更新:基于客户旧画像(可能没有)+ 聊天记录重新生成全量动态栏目。
// clientId 优先;否则按 wxid(托管客户)定位/自动建档。忽略 analyzed 标记、不回写,与定时迭代互不干扰。
func (a *AdviceIterateBiz) RegenerateProfile(ctx context.Context, param *entitys.AdvicerProfileRegenerateReq) (res *entitys.AdvicerIterateRunRes, err error) {
res = &entitys.AdvicerIterateRunRes{Details: []string{}}
if param == nil || (len(param.ClientId) == 0 && len(param.Wxid) == 0) {
return res, errors.New("请指定客户(clientId 或 wxid)")
}
var c *mongo_model.AdvicerClientItem
if len(param.ClientId) != 0 {
info, e := a.clientBiz.Info(ctx, &entitys.AdvicerClientInfoReq{Id: param.ClientId})
if e != nil {
return res, fmt.Errorf("客户不存在: %w", e)
}
objID, e := primitive.ObjectIDFromHex(param.ClientId)
if e != nil {
return res, e
}
c = &mongo_model.AdvicerClientItem{Id: objID, AdvicerClientMongo: info}
} else {
advicerId := int32(0)
if len(param.AppId) != 0 {
if adv, e := a.advicerBiz.FindByWxDeviceId(ctx, param.AppId); e == nil {
advicerId = adv.AdvicerID
}
}
item, _, e := a.clientBiz.EnsureByWxid(ctx, param.Wxid, param.ProjectId, advicerId, param.AppId)
if e != nil {
return res, fmt.Errorf("客户解析失败: %w", e)
}
c = item
}
if len(c.Wxid) == 0 {
return res, errors.New("该客户未绑定微信 wxid,无聊天记录可生成画像")
}
msgs := a.recentMsgs(ctx, c.Wxid, 200)
if len(msgs) == 0 {
res.Details = append(res.Details, clientNameOf(c)+":暂无聊天记录,未生成画像")
return res, nil
}
res.Scanned = len(msgs)
_, model := projectModelOf(ctx, a.projectImpl, a.modelSupImpl, c.ProjectId)
if model == nil {
return res, errors.New("项目未配置可用模型(请检查 modelSupId)")
}
manual, wxRemark := a.fetchRemarks(ctx, c)
content := buildRegenerateUserContent(c, msgs, manual, wxRemark)
messages := []openai.ResponseInputMessage{
{Role: openai.ChatMessageRoleSystem, Content: constants.ClientIteratePrompt},
{Role: openai.ChatMessageRoleUser, Content: content},
}
resp, err := a.openai.CreateResponseMessages(ctx, model.Key, model.URL, model.ChatModel, messages, "")
if err != nil {
return res, fmt.Errorf("分析聊天记录失败: %w", err)
}
var patch iteratePatch
if err = json.Unmarshal([]byte(extractJsonObject(resp.GetOutputText())), &patch); err != nil {
return res, fmt.Errorf("解析画像结果失败: %w", err)
}
updated, err := a.applyIterateSections(ctx, c, patch.Sections, true)
if err != nil {
return res, err
}
if updated {
res.Updated++
text := clientNameOf(c) + ":画像已更新"
if s := strings.TrimSpace(patch.Summary); s != "" {
text += "(" + s + ")"
}
res.Details = append(res.Details, text)
} else {
res.Details = append(res.Details, clientNameOf(c)+":未能从聊天记录中提取到画像信息")
}
return res, nil
}
// clientNameOf 客户显示名(画像「姓名」栏目优先,其次截断 wxid),供迭代/生成结果描述使用。
func clientNameOf(c *mongo_model.AdvicerClientItem) string {
if n := strings.TrimSpace(c.SectionString(mongo_model.SectionKeyClientName)); n != "" {
return n
}
if len(c.Wxid) > 10 {
return c.Wxid[:10] + "..."
}
return c.Wxid
}
// buildRegenerateUserContent 组装画像更新输入(现有画像 + 客户备注 + 全量聊天记录)
func buildRegenerateUserContent(c *mongo_model.AdvicerClientItem, msgs []mongo_model.AdvicerWxMsgMongo, manual, wxRemark string) string {
var b strings.Builder
b.WriteString("[客户现有画像]\n")
b.WriteString(pkg.JsonStringIgonErr(c.Entity()))
if manual != "" {
b.WriteString("\n[人工备注]" + manual)
}
if wxRemark != "" {
b.WriteString("\n[微信备注名]" + wxRemark)
}
b.WriteString("\n\n[需要分析的聊天记录(按时间先后)]\n")
for i := range msgs {
b.WriteString(formatIterateMsg(&msgs[i]))
}
return b.String()
}
// buildIterateUserContent 组装迭代输入(现有画像 + 客户备注 + 本轮新增对话 + 更早上下文)
func buildIterateUserContent(c *mongo_model.AdvicerClientItem, history []mongo_model.AdvicerWxMsgMongo, newMsgs []mongo_model.AdvicerWxMsgItem, manual, wxRemark string) string {
var b strings.Builder
b.WriteString("[客户现有画像]\n")
b.WriteString(pkg.JsonStringIgonErr(c.Entity()))
if manual != "" {
b.WriteString("\n[人工备注]" + manual)
}
if wxRemark != "" {
b.WriteString("\n[微信备注名]" + wxRemark)
}
b.WriteString("\n\n[本轮需要分析的聊天记录(按时间先后)]\n")
if len(newMsgs) == 0 {
b.WriteString("(无)\n")
}
for i := range newMsgs {
b.WriteString(formatIterateMsg(&newMsgs[i].AdvicerWxMsgMongo))
}
// 更早的上下文(仅帮助理解语境)
if len(newMsgs) > 0 {
earliest := newMsgs[0].CreateAt
var ctxBuilder strings.Builder
for i := range history {
if history[i].CreateAt.Before(earliest) {
ctxBuilder.WriteString(formatIterateMsg(&history[i]))
}
}
if ctxBuilder.Len() > 0 {
b.WriteString("\n[更早的聊天记录(仅供理解上下文,无需重复提取)]\n")
b.WriteString(ctxBuilder.String())
}
}
return b.String()
}
// formatIterateMsg 单条消息格式化(带时间)
func formatIterateMsg(m *mongo_model.AdvicerWxMsgMongo) string {
who := "客户"
if m.Direction != mongo_model.WxMsgDirectionCustomer {
who = "我"
}
content := m.Content
if m.MsgType != mongo_model.WxMsgTypeText {
content = "[" + m.MsgType + "消息] " + content
}
return fmt.Sprintf("[%s] %s:%s\n", m.CreateAt.Format("01-02 15:04"), who, content)
}
// mergeUnique 追加去重(保留原有顺序)
func mergeUnique(old []string, add []string) []string {
seen := map[string]bool{}
for _, s := range old {
seen[strings.TrimSpace(s)] = true
}
merged := append([]string{}, old...)
for _, s := range add {
s = strings.TrimSpace(s)
if s == "" || seen[s] {
continue
}
seen[s] = true
merged = append(merged, s)
}
return merged
}