package biz import ( "ai_scheduler/internal/config" "ai_scheduler/internal/data/constants" "ai_scheduler/internal/data/mongo_model" "ai_scheduler/internal/entitys" "ai_scheduler/internal/pkg" "context" "errors" "fmt" "strings" "time" "go.mongodb.org/mongo-driver/bson" "go.mongodb.org/mongo-driver/mongo/options" ) // 触达结果状态 const ( proactiveNone = iota // 未适用(非生日/无活动可推) proactiveSent // 已发送 proactiveSkipped // 已跳过(去重/频控等) proactiveFailed // 发送失败 ) // AdviceProactiveBiz AI 主动对话(生日祝福 / 新活动推送) // 规则:熟客与非客户不主动触达;意向客户可生日+活动触达;沉睡客户只推活动; // 同一客户同一活动/同一年度生日只触达一次,且默认每类型每天最多一次;自动链路受活跃时段与总开关限制。 type AdviceProactiveBiz struct { cfg *config.Config mongo *pkg.Mongo clientBiz *AdviceClientBiz proactiveMongo *mongo_model.AdvicerProactiveLogMongo activityBiz *AdviceActivityBiz strategyBiz *AdviceStrategyBiz sendBiz *AdviceWxSendBiz projectBiz *AdviceProjectBiz } func NewAdviceProactiveBiz( cfg *config.Config, mongo *pkg.Mongo, clientBiz *AdviceClientBiz, proactiveMongo *mongo_model.AdvicerProactiveLogMongo, activityBiz *AdviceActivityBiz, strategyBiz *AdviceStrategyBiz, sendBiz *AdviceWxSendBiz, projectBiz *AdviceProjectBiz, ) *AdviceProactiveBiz { return &AdviceProactiveBiz{ cfg: cfg, mongo: mongo, clientBiz: clientBiz, proactiveMongo: proactiveMongo, activityBiz: activityBiz, strategyBiz: strategyBiz, sendBiz: sendBiz, projectBiz: projectBiz, } } // getProjectWxToken 获取项目级 wx_token,优先项目配置,回退全局配置 func (a *AdviceProactiveBiz) getProjectWxToken(ctx context.Context, projectId int32) string { if a.projectBiz != nil && projectId != 0 { if info, err := a.projectBiz.BaseInfo(projectId); err == nil && info.WxToken != "" { return info.WxToken } } return a.cfg.Advicer.WxToken } // Scan 主动触达扫描(生日祝福 + 新活动推送),定时任务每小时调用,也可后台手动触发。 // 自动链路:受 AutoReply 总开关与活跃时段(active_hour_start/end)限制;Force=true 忽略两者。 func (a *AdviceProactiveBiz) Scan(ctx context.Context, param *entitys.AdvicerProactiveRunReq) (res *entitys.AdvicerProactiveRunRes, err error) { res = &entitys.AdvicerProactiveRunRes{Details: []string{}} if param == nil { param = &entitys.AdvicerProactiveRunReq{} } if !param.Force { if !a.cfg.Advicer.AutoReply { return res, nil } if !a.inActiveWindow() { return res, nil } } // wx_token 从项目级获取,此处不再做全局 Available 检查 filter := bson.M{"aiEnabled": true, "wxid": bson.M{"$ne": ""}} if param.ProjectId != 0 { filter["projectId"] = param.ProjectId } if param.AdvicerId != 0 { filter["advicerId"] = param.AdvicerId } cursor, err := a.mongo.Co(a.clientBiz.AdvicerClientMongo).Find(ctx, filter) if err != nil { return res, err } var clients []mongo_model.AdvicerClientItem for cursor.Next(ctx) { var item mongo_model.AdvicerClientItem if e := cursor.Decode(&item); e != nil { return res, e } clients = append(clients, item) } if err = cursor.Err(); err != nil { return res, err } res.Scanned = len(clients) todayMD := time.Now().Format("01-02") for i := range clients { item := clients[i] // 熟客:AI 不介入(人工服务);非客户:不参与 if item.ClientLevel == mongo_model.ClientLevelRegular || item.ClientLevel == mongo_model.ClientLevelNon { res.Skipped++ continue } // 1) 生日触达 if normalizeBirthday(item.PersonalInfo.Birthday) == todayMD { stat, detail := a.sendBirthday(ctx, &item, param, a.getProjectWxToken(ctx, item.ProjectId)) if detail != "" { res.Details = append(res.Details, detail) } switch stat { case proactiveSent: res.Birthday++ case proactiveFailed: res.Failed++ case proactiveSkipped: res.Skipped++ } } // 2) 活动触达 stat, detail := a.sendActivityForClient(ctx, &item, param, a.getProjectWxToken(ctx, item.ProjectId)) if detail != "" { res.Details = append(res.Details, detail) } switch stat { case proactiveSent: res.Activity++ case proactiveFailed: res.Failed++ case proactiveSkipped: res.Skipped++ } } return res, nil } // sendBirthday 生日祝福触达(同一客户每年度一次) func (a *AdviceProactiveBiz) sendBirthday(ctx context.Context, c *mongo_model.AdvicerClientItem, param *entitys.AdvicerProactiveRunReq, wxToken string) (stat int, detail string) { clientId := c.Id.Hex() refId := time.Now().Format("2006") if a.hasLog(ctx, clientId, mongo_model.ProactiveTypeBirthday, refId) { return proactiveSkipped, clientName(c) + ":本年度生日祝福已触达过" } if !param.Force && a.hasLogToday(ctx, clientId, mongo_model.ProactiveTypeBirthday) { return proactiveSkipped, clientName(c) + ":今日已触达过" } if param.DryRun { return proactiveSkipped, clientName(c) + ":今日生日,待发送祝福(dry-run)" } if len(wxToken) == 0 { return proactiveFailed, clientName(c) + ":项目未配置 wx_token,无法发送消息" } replies, err := a.strategyBiz.GenerateReply(ctx, c, constants.StrategyBirthdayMission, "", "") if err != nil || len(replies) == 0 { a.logContact(ctx, c, mongo_model.ProactiveTypeBirthday, refId, "", err) return proactiveFailed, fmt.Sprintf("%s:生日祝福生成失败: %v", clientName(c), err) } sentList, e := a.sendBiz.SendMultiWithToken(ctx, wxToken, c.AppId, c.Wxid, replies, mongo_model.WxMsgSourceWeb) a.logContact(ctx, c, mongo_model.ProactiveTypeBirthday, refId, strings.Join(sentList, "\n"), e) if e != nil { return proactiveFailed, fmt.Sprintf("%s:生日祝福发送失败: %v", clientName(c), e) } return proactiveSent, clientName(c) + ":已发送生日祝福" } // sendActivityForClient 活动推送触达(同一客户同一活动一次,每类型每天一次) func (a *AdviceProactiveBiz) sendActivityForClient(ctx context.Context, c *mongo_model.AdvicerClientItem, param *entitys.AdvicerProactiveRunReq, wxToken string) (stat int, detail string) { acts, _ := a.activityBiz.ActiveList(ctx, c.ProjectId, c.ClientLevel) if len(acts) == 0 { return proactiveNone, "" } clientId := c.Id.Hex() for _, act := range acts { refId := act.Id.Hex() if a.hasLog(ctx, clientId, mongo_model.ProactiveTypeActivity, refId) { // 该活动已推送过,尝试下一个未推送的活动 continue } if !param.Force && a.hasLogToday(ctx, clientId, mongo_model.ProactiveTypeActivity) { return proactiveSkipped, clientName(c) + ":今日已推送过活动" } if param.DryRun { return proactiveSkipped, fmt.Sprintf("%s:待推送活动「%s」(dry-run)", clientName(c), act.Name) } return a.pushActivity(ctx, c, act.AdvicerActivityMongo, refId, wxToken) } return proactiveNone, "" } // pushActivity 生成活动推送内容并发送(单客户),返回结果状态 func (a *AdviceProactiveBiz) pushActivity(ctx context.Context, c *mongo_model.AdvicerClientItem, act mongo_model.AdvicerActivityMongo, refId string, wxToken string) (stat int, detail string) { if len(wxToken) == 0 { return proactiveFailed, clientName(c) + ":项目未配置 wx_token,无法发送消息" } replies, err := a.strategyBiz.GenerateReply(ctx, c, constants.StrategyProactiveActivityMission, buildActivityCtx(act), "") if err != nil || len(replies) == 0 { a.logContact(ctx, c, mongo_model.ProactiveTypeActivity, refId, "", err) return proactiveFailed, fmt.Sprintf("%s:活动推送生成失败: %v", clientName(c), err) } sentList, e := a.sendBiz.SendMultiWithToken(ctx, wxToken, c.AppId, c.Wxid, replies, mongo_model.WxMsgSourceWeb) a.logContact(ctx, c, mongo_model.ProactiveTypeActivity, refId, strings.Join(sentList, "\n"), e) if e != nil { return proactiveFailed, fmt.Sprintf("%s:活动推送发送失败: %v", clientName(c), e) } return proactiveSent, fmt.Sprintf("%s:已推送活动「%s」", clientName(c), act.Name) } // PushActivity 将指定活动立即推送给目标客户群(后台活动列表"推送"按钮;跳过熟客/非客户与已推送客户) func (a *AdviceProactiveBiz) PushActivity(ctx context.Context, param *entitys.AdvicerActivityPushReq) (res *entitys.AdvicerProactiveRunRes, err error) { res = &entitys.AdvicerProactiveRunRes{Details: []string{}} if param == nil || len(param.Id) == 0 { return res, errors.New("活动ID不能为空") } act, err := a.activityBiz.Info(ctx, &entitys.AdvicerActivityInfoReq{Id: param.Id}) if err != nil { return res, fmt.Errorf("活动不存在: %w", err) } // wx_token 从项目级获取,不再做全局 Available 检查 filter := bson.M{"aiEnabled": true, "wxid": bson.M{"$ne": ""}} if act.ProjectId != 0 { filter["projectId"] = act.ProjectId } if act.AdvicerId != 0 { filter["advicerId"] = act.AdvicerId } cursor, err := a.mongo.Co(a.clientBiz.AdvicerClientMongo).Find(ctx, filter) if err != nil { return res, err } var clients []mongo_model.AdvicerClientItem for cursor.Next(ctx) { var item mongo_model.AdvicerClientItem if e := cursor.Decode(&item); e != nil { return res, e } clients = append(clients, item) } if err = cursor.Err(); err != nil { return res, err } res.Scanned = len(clients) for i := range clients { item := clients[i] if item.ClientLevel == mongo_model.ClientLevelRegular || item.ClientLevel == mongo_model.ClientLevelNon { res.Skipped++ continue } // 活动目标等级过滤(空=全部) if len(act.TargetLevels) > 0 && !containsStr(act.TargetLevels, item.ClientLevel) { res.Skipped++ continue } // 该活动已推送过 if a.hasLog(ctx, item.Id.Hex(), mongo_model.ProactiveTypeActivity, param.Id) { res.Skipped++ continue } if param.DryRun { res.Skipped++ res.Details = append(res.Details, fmt.Sprintf("%s:待推送活动「%s」(dry-run)", clientName(&item), act.Name)) continue } wxToken := a.getProjectWxToken(ctx, item.ProjectId) stat, detail := a.pushActivity(ctx, &item, act, param.Id, wxToken) if detail != "" { res.Details = append(res.Details, detail) } switch stat { case proactiveSent: res.Activity++ case proactiveFailed: res.Failed++ default: res.Skipped++ } } return res, nil } // SendToClient 手动推送给指定客户(活动/自定义话术),后台手动触发不受总开关与时段限制 func (a *AdviceProactiveBiz) SendToClient(ctx context.Context, param *entitys.AdvicerProactiveSendReq) (res *entitys.AdvicerProactiveRunRes, err error) { res = &entitys.AdvicerProactiveRunRes{Details: []string{}} if param == nil || len(param.ClientId) == 0 { return res, errors.New("客户ID不能为空") } client, err := a.clientBiz.Info(ctx, &entitys.AdvicerClientInfoReq{Id: param.ClientId}) if err != nil { return res, fmt.Errorf("客户不存在: %w", err) } if client.ClientLevel == mongo_model.ClientLevelRegular { return res, errors.New("熟客由人工服务,不建议 AI 主动触达") } if client.ClientLevel == mongo_model.ClientLevelNon { return res, errors.New("非客户不参与 AI 触达") } if !client.AiEnabled { return res, errors.New("客户 AI 参与开关已关闭") } if len(client.Wxid) == 0 { return res, errors.New("客户未绑定微信 wxid") } // wx_token 从项目级获取,不再做全局 Available 检查 item := &mongo_model.AdvicerClientItem{AdvicerClientMongo: client} var replies []string var refId string switch { case strings.TrimSpace(param.Content) != "": // 自定义话术:直接拆条发送,不去重(手动行为) replies = splitReplies(param.Content) refId = "manual-" + time.Now().Format("20060102150405") if len(replies) == 0 { return res, errors.New("自定义内容为空") } case len(param.ActivityId) != 0: act, e := a.activityBiz.Info(ctx, &entitys.AdvicerActivityInfoReq{Id: param.ActivityId}) if e != nil { return res, fmt.Errorf("活动不存在: %w", e) } refId = param.ActivityId replies, err = a.strategyBiz.GenerateReply(ctx, item, constants.StrategyProactiveActivityMission, buildActivityCtx(act), "") default: // 未指定活动:取一个当前生效活动作为素材,无则退化为日常跟进 acts, _ := a.activityBiz.ActiveList(ctx, client.ProjectId, client.ClientLevel) mission := constants.StrategyFollowupMission extra := "" refId = "manual-" + time.Now().Format("20060102150405") if len(acts) > 0 { mission = constants.StrategyProactiveActivityMission extra = buildActivityCtx(acts[0].AdvicerActivityMongo) refId = acts[0].Id.Hex() } replies, err = a.strategyBiz.GenerateReply(ctx, item, mission, extra, "") } if err != nil { return res, err } if len(replies) == 0 { return res, errors.New("未生成可发送内容") } wxToken := a.getProjectWxToken(ctx, client.ProjectId) if len(wxToken) == 0 { return res, errors.New("项目未配置 wx_token,无法发送消息") } sentList, e := a.sendBiz.SendMultiWithToken(ctx, wxToken, client.AppId, client.Wxid, replies, mongo_model.WxMsgSourceWeb) a.logContactByWxid(ctx, client, refId, strings.Join(sentList, "\n"), e) if e != nil { return res, e } res.Activity = 1 res.Details = append(res.Details, "已发送:\n"+strings.Join(sentList, "\n")) return res, nil } // LogList 主动触达记录查询(按时间倒序) func (a *AdviceProactiveBiz) LogList(ctx context.Context, param *entitys.AdvicerProactiveLogListReq) (list []mongo_model.AdvicerProactiveLogItem, err error) { filter := bson.M{} if param.ClientId != "" { filter["clientId"] = param.ClientId } if param.Wxid != "" { filter["wxid"] = param.Wxid } if param.Type != "" { filter["type"] = param.Type } opts := options.Find().SetSort(bson.D{{Key: "sendAt", 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.proactiveMongo).Find(ctx, filter, opts) if err != nil { return nil, err } for cursor.Next(ctx) { var item mongo_model.AdvicerProactiveLogItem 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 } // ==================== 内部工具 ==================== // inActiveWindow 是否处于主动触达活跃时段(默认 9-21 点,可配置,支持跨天) func (a *AdviceProactiveBiz) inActiveWindow() bool { start := a.cfg.Advicer.ActiveHourStart end := a.cfg.Advicer.ActiveHourEnd if start <= 0 || start > 23 { start = 9 } if end <= 0 || end > 23 { end = 21 } h := time.Now().Hour() if start <= end { return h >= start && h < end } return h >= start || h < end } // hasLog 是否已有成功触达记录(同一客户+类型+去重键) func (a *AdviceProactiveBiz) hasLog(ctx context.Context, clientId, typ, refId string) bool { count, err := a.mongo.Co(a.proactiveMongo).CountDocuments(ctx, bson.M{ "clientId": clientId, "type": typ, "refId": refId, "status": mongo_model.ProactiveStatusSent, }) return err == nil && count > 0 } // hasLogToday 当日是否已触达过该类型(频控) func (a *AdviceProactiveBiz) hasLogToday(ctx context.Context, clientId, typ string) bool { start, e := time.ParseInLocation("2006-01-02", time.Now().Format("2006-01-02"), time.Local) if e != nil { start = time.Now().Add(-24 * time.Hour) } count, err := a.mongo.Co(a.proactiveMongo).CountDocuments(ctx, bson.M{ "clientId": clientId, "type": typ, "status": mongo_model.ProactiveStatusSent, "sendAt": bson.M{"$gte": start}, }) return err == nil && count > 0 } // logContact 记录触达日志(err 非空记为失败) func (a *AdviceProactiveBiz) logContact(ctx context.Context, c *mongo_model.AdvicerClientItem, typ, refId, content string, err error) { now := time.Now() log := &mongo_model.AdvicerProactiveLogMongo{ ProjectId: c.ProjectId, AdvicerId: c.AdvicerId, ClientId: c.Id.Hex(), Wxid: c.Wxid, Type: typ, RefId: refId, Content: content, Status: mongo_model.ProactiveStatusSent, SendAt: now, CreateAt: now, } if err != nil { log.Status = mongo_model.ProactiveStatusFailed log.ErrMsg = err.Error() } _, _ = a.mongo.Co(a.proactiveMongo).InsertOne(ctx, log) } // logContactByWxid 记录触达日志(无客户 _id 场景,手动推送) func (a *AdviceProactiveBiz) logContactByWxid(ctx context.Context, c mongo_model.AdvicerClientMongo, refId, content string, err error) { now := time.Now() log := &mongo_model.AdvicerProactiveLogMongo{ ProjectId: c.ProjectId, AdvicerId: c.AdvicerId, Wxid: c.Wxid, Type: mongo_model.ProactiveTypeActivity, RefId: refId, Content: content, Status: mongo_model.ProactiveStatusSent, SendAt: now, CreateAt: now, } if err != nil { log.Status = mongo_model.ProactiveStatusFailed log.ErrMsg = err.Error() } _, _ = a.mongo.Co(a.proactiveMongo).InsertOne(ctx, log) } // normalizeBirthday 归一化生日为 MM-DD(支持 "05-20"/"1995-05-20"/"5-20"/"05/20"),无法识别返回空串 func normalizeBirthday(b string) string { b = strings.TrimSpace(b) if b == "" { return "" } if len(b) >= 10 { b = b[5:10] // "1995-05-20" → "05-20" } b = strings.ReplaceAll(b, "/", "-") parts := strings.Split(b, "-") if len(parts) != 2 { return "" } m := strings.TrimLeft(parts[0], "0") d := strings.TrimLeft(parts[1], "0") if len(m) == 0 || len(d) == 0 || len(m) > 2 || len(d) > 2 { return "" } if len(m) == 1 { m = "0" + m } if len(d) == 1 { d = "0" + d } return m + "-" + d } // clientName 客户显示名(姓名优先,其次截断的 wxid) func clientName(c *mongo_model.AdvicerClientItem) string { if n := strings.TrimSpace(c.PersonalInfo.Name); n != "" { return n } if len(c.Wxid) > 10 { return c.Wxid[:10] + "..." } return c.Wxid } // containsStr 字符串切片包含判断 func containsStr(list []string, v string) bool { for _, item := range list { if item == v { return true } } return false }