You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
Open-IM-Server/pkg/common/db/controller/msg.go

729 lines
29 KiB

2 years ago
package controller
2 years ago
import (
2 years ago
"OpenIM/pkg/common/config"
2 years ago
"OpenIM/pkg/common/constant"
"OpenIM/pkg/common/db/cache"
unRelationTb "OpenIM/pkg/common/db/table/unrelation"
"OpenIM/pkg/common/db/unrelation"
2 years ago
"OpenIM/pkg/common/kafka"
2 years ago
"OpenIM/pkg/common/log"
"OpenIM/pkg/common/prome"
"OpenIM/pkg/common/tracelog"
2 years ago
"fmt"
2 years ago
"github.com/gogo/protobuf/sortkeys"
"sync"
2 years ago
"time"
2 years ago
2 years ago
pbMsg "OpenIM/pkg/proto/msg"
"OpenIM/pkg/proto/sdkws"
"OpenIM/pkg/utils"
2 years ago
"context"
2 years ago
"errors"
"github.com/go-redis/redis/v8"
"go.mongodb.org/mongo-driver/mongo"
"github.com/golang/protobuf/proto"
2 years ago
)
2 years ago
type MsgDatabase interface {
2 years ago
// 批量插入消息
2 years ago
BatchInsertChat2DB(ctx context.Context, sourceID string, msgList []*pbMsg.MsgDataToMQ, currentMaxSeq int64) error
2 years ago
// 刪除redis中消息缓存
2 years ago
DeleteMessageFromCache(ctx context.Context, sourceID string, msgList []*pbMsg.MsgDataToMQ) error
2 years ago
// incrSeq然后批量插入缓存
2 years ago
BatchInsertChat2Cache(ctx context.Context, sourceID string, msgList []*pbMsg.MsgDataToMQ) (int64, error)
2 years ago
// 删除消息 返回不存在的seqList
2 years ago
DelMsgBySeqs(ctx context.Context, userID string, seqs []int64) (totalUnExistSeqs []int64, err error)
2 years ago
// 获取群ID或者UserID最新一条在mongo里面的消息
2 years ago
// 通过seqList获取mongo中写扩散消息
2 years ago
GetMsgBySeqs(ctx context.Context, userID string, seqs []int64) (seqMsg []*sdkws.MsgData, err error)
2 years ago
// 通过seqList获取大群在 mongo里面的消息
2 years ago
GetSuperGroupMsgBySeqs(ctx context.Context, groupID string, seqs []int64) (seqMsg []*sdkws.MsgData, err error)
2 years ago
// 删除用户所有消息/redis/mongo然后重置seq
2 years ago
CleanUpUserMsg(ctx context.Context, userID string) error
2 years ago
// 删除大群消息重置群成员最小群seq, remainTime为消息保留的时间单位秒,超时消息删除, 传0删除所有消息(此方法不删除 redis cache)
2 years ago
DeleteUserSuperGroupMsgsAndSetMinSeq(ctx context.Context, groupID string, userIDs []string, remainTime int64) error
2 years ago
// 删除用户消息重置最小seq remainTime为消息保留的时间单位秒,超时消息删除, 传0删除所有消息(此方法不删除redis cache)
DeleteUserMsgsAndSetMinSeq(ctx context.Context, userID string, remainTime int64) error
2 years ago
// 获取用户 seq mongo和redis
GetUserMinMaxSeqInMongoAndCache(ctx context.Context, userID string) (minSeqMongo, maxSeqMongo, minSeqCache, maxSeqCache int64, err error)
// 获取群 seq mongo和redis
GetSuperGroupMinMaxSeqInMongoAndCache(ctx context.Context, groupID string) (minSeqMongo, maxSeqMongo, maxSeqCache int64, err error)
// 设置群用户最小seq 直接调用cache
SetGroupUserMinSeq(ctx context.Context, groupID, userID string, minSeq int64) (err error)
2 years ago
GetGroupUserMinSeq(ctx context.Context, groupID, userID string) (int64, error)
2 years ago
// 设置用户最小seq 直接调用cache
SetUserMinSeq(ctx context.Context, userID string, minSeq int64) (err error)
2 years ago
2 years ago
JudgeMessageReactionExist(ctx context.Context, clientMsgID string, sessionType int32) (bool, error)
2 years ago
SetMessageTypeKeyValue(ctx context.Context, clientMsgID string, sessionType int32, typeKey, value string) error
SetMessageReactionExpire(ctx context.Context, clientMsgID string, sessionType int32, expiration time.Duration) (bool, error)
GetExtendMsg(ctx context.Context, sourceID string, sessionType int32, clientMsgID string, maxMsgUpdateTime int64) (*pbMsg.ExtendMsg, error)
2 years ago
InsertOrUpdateReactionExtendMsgSet(ctx context.Context, sourceID string, sessionType int32, clientMsgID string, msgFirstModifyTime int64, reactionExtensionList map[string]*sdkws.KeyValue) error
GetMessageTypeKeyValue(ctx context.Context, clientMsgID string, sessionType int32, typeKey string) (string, error)
GetOneMessageAllReactionList(ctx context.Context, clientMsgID string, sessionType int32) (map[string]string, error)
DeleteOneMessageKey(ctx context.Context, clientMsgID string, sessionType int32, subKey string) error
DeleteReactionExtendMsgSet(ctx context.Context, sourceID string, sessionType int32, clientMsgID string, msgFirstModifyTime int64, reactionExtensionList map[string]*sdkws.KeyValue) error
2 years ago
SetSendMsgStatus(ctx context.Context, id string, status int32) error
GetSendMsgStatus(ctx context.Context, id string) (int32, error)
2 years ago
GetUserMaxSeq(ctx context.Context, userID string) (int64, error)
GetUserMinSeq(ctx context.Context, userID string) (int64, error)
GetGroupMaxSeq(ctx context.Context, groupID string) (int64, error)
GetGroupMinSeq(ctx context.Context, groupID string) (int64, error)
2 years ago
MsgToMQ(ctx context.Context, key string, msg2mq *pbMsg.MsgDataToMQ) error
MsgToModifyMQ(ctx context.Context, aggregationID string, triggerID string, messages []*pbMsg.MsgDataToMQ) error
MsgToPushMQ(ctx context.Context, sourceID string, msg2mq *pbMsg.MsgDataToMQ) error
MsgToMongoMQ(ctx context.Context, aggregationID string, triggerID string, messages []*pbMsg.MsgDataToMQ, lastSeq int64) error
2 years ago
}
2 years ago
2 years ago
func NewMsgDatabase(msgDocModel unRelationTb.MsgDocModelInterface, cacheModel cache.Model) MsgDatabase {
return &msgDatabase{
2 years ago
msgDocDatabase: msgDocModel,
cache: cacheModel,
producer: kafka.NewKafkaProducer(config.Config.Kafka.Ws2mschat.Addr, config.Config.Kafka.Ws2mschat.Topic),
producerToMongo: kafka.NewKafkaProducer(config.Config.Kafka.MsgToMongo.Addr, config.Config.Kafka.MsgToMongo.Topic),
producerToPush: kafka.NewKafkaProducer(config.Config.Kafka.Ms2pschat.Addr, config.Config.Kafka.Ms2pschat.Topic),
producerToModify: kafka.NewKafkaProducer(config.Config.Kafka.MsgToModify.Addr, config.Config.Kafka.MsgToModify.Topic),
2 years ago
}
}
func InitMsgDatabase(rdb redis.UniversalClient, database *mongo.Database) MsgDatabase {
cacheModel := cache.NewCacheModel(rdb)
msgDocModel := unrelation.NewMsgMongoDriver(database)
msgDatabase := NewMsgDatabase(msgDocModel, cacheModel)
return msgDatabase
2 years ago
}
2 years ago
type msgDatabase struct {
2 years ago
msgDocDatabase unRelationTb.MsgDocModelInterface
extendMsgDatabase unRelationTb.ExtendMsgSetModelInterface
cache cache.Model
2 years ago
producer *kafka.Producer
2 years ago
producerToMongo *kafka.Producer
producerToModify *kafka.Producer
producerToPush *kafka.Producer
2 years ago
// model
2 years ago
msg unRelationTb.MsgDocModel
extendMsgSetModel unRelationTb.ExtendMsgSetModel
2 years ago
}
2 years ago
func (db *msgDatabase) JudgeMessageReactionExist(ctx context.Context, clientMsgID string, sessionType int32) (bool, error) {
return db.cache.JudgeMessageReactionExist(ctx, clientMsgID, sessionType)
2 years ago
}
2 years ago
func (db *msgDatabase) SetMessageTypeKeyValue(ctx context.Context, clientMsgID string, sessionType int32, typeKey, value string) error {
2 years ago
return db.cache.SetMessageTypeKeyValue(ctx, clientMsgID, sessionType, typeKey, value)
2 years ago
}
2 years ago
func (db *msgDatabase) SetMessageReactionExpire(ctx context.Context, clientMsgID string, sessionType int32, expiration time.Duration) (bool, error) {
2 years ago
return db.cache.SetMessageReactionExpire(ctx, clientMsgID, sessionType, expiration)
2 years ago
}
2 years ago
func (db *msgDatabase) GetMessageTypeKeyValue(ctx context.Context, clientMsgID string, sessionType int32, typeKey string) (string, error) {
2 years ago
return db.cache.GetMessageTypeKeyValue(ctx, clientMsgID, sessionType, typeKey)
2 years ago
}
2 years ago
func (db *msgDatabase) GetOneMessageAllReactionList(ctx context.Context, clientMsgID string, sessionType int32) (map[string]string, error) {
2 years ago
return db.cache.GetOneMessageAllReactionList(ctx, clientMsgID, sessionType)
2 years ago
}
2 years ago
func (db *msgDatabase) DeleteOneMessageKey(ctx context.Context, clientMsgID string, sessionType int32, subKey string) error {
2 years ago
return db.cache.DeleteOneMessageKey(ctx, clientMsgID, sessionType, subKey)
2 years ago
}
2 years ago
func (db *msgDatabase) InsertOrUpdateReactionExtendMsgSet(ctx context.Context, sourceID string, sessionType int32, clientMsgID string, msgFirstModifyTime int64, reactionExtensions map[string]*sdkws.KeyValue) error {
2 years ago
return db.extendMsgDatabase.InsertOrUpdateReactionExtendMsgSet(ctx, sourceID, sessionType, clientMsgID, msgFirstModifyTime, db.extendMsgSetModel.Pb2Model(reactionExtensions))
2 years ago
}
2 years ago
func (db *msgDatabase) GetExtendMsg(ctx context.Context, sourceID string, sessionType int32, clientMsgID string, maxMsgUpdateTime int64) (*pbMsg.ExtendMsg, error) {
2 years ago
extendMsgSet, err := db.extendMsgDatabase.GetExtendMsgSet(ctx, sourceID, sessionType, maxMsgUpdateTime)
2 years ago
if err != nil {
return nil, err
}
extendMsg, ok := extendMsgSet.ExtendMsgs[clientMsgID]
if !ok {
return nil, errors.New(fmt.Sprintf("cant find client msg id: %s", clientMsgID))
}
reactionExtensionList := make(map[string]*pbMsg.KeyValueResp)
for key, model := range extendMsg.ReactionExtensionList {
reactionExtensionList[key] = &pbMsg.KeyValueResp{
KeyValue: &sdkws.KeyValue{
TypeKey: model.TypeKey,
Value: model.Value,
LatestUpdateTime: model.LatestUpdateTime,
},
}
}
return &pbMsg.ExtendMsg{
2 years ago
ReactionExtensions: reactionExtensionList,
ClientMsgID: extendMsg.ClientMsgID,
MsgFirstModifyTime: extendMsg.MsgFirstModifyTime,
AttachedInfo: extendMsg.AttachedInfo,
Ex: extendMsg.Ex,
2 years ago
}, nil
2 years ago
}
2 years ago
func (db *msgDatabase) DeleteReactionExtendMsgSet(ctx context.Context, sourceID string, sessionType int32, clientMsgID string, msgFirstModifyTime int64, reactionExtensions map[string]*sdkws.KeyValue) error {
2 years ago
return db.extendMsgDatabase.DeleteReactionExtendMsgSet(ctx, sourceID, sessionType, clientMsgID, msgFirstModifyTime, db.extendMsgSetModel.Pb2Model(reactionExtensions))
2 years ago
}
2 years ago
func (db *msgDatabase) SetSendMsgStatus(ctx context.Context, id string, status int32) error {
2 years ago
return db.cache.SetSendMsgStatus(ctx, id, status)
2 years ago
}
2 years ago
func (db *msgDatabase) GetSendMsgStatus(ctx context.Context, id string) (int32, error) {
2 years ago
return db.cache.GetSendMsgStatus(ctx, id)
2 years ago
}
2 years ago
func (db *msgDatabase) MsgToMQ(ctx context.Context, key string, msg2mq *pbMsg.MsgDataToMQ) error {
_, _, err := db.producer.SendMessage(ctx, key, msg2mq)
return err
2 years ago
}
2 years ago
func (db *msgDatabase) MsgToModifyMQ(ctx context.Context, aggregationID string, triggerID string, messages []*pbMsg.MsgDataToMQ) error {
if len(messages) > 0 {
_, _, err := db.producerToModify.SendMessage(ctx, aggregationID, &pbMsg.MsgDataToModifyByMQ{AggregationID: aggregationID, Messages: messages, TriggerID: triggerID})
return err
}
return nil
}
func (db *msgDatabase) MsgToPushMQ(ctx context.Context, key string, msg2mq *pbMsg.MsgDataToMQ) error {
mqPushMsg := pbMsg.PushMsgDataToMQ{MsgData: msg2mq.MsgData, SourceID: key}
_, _, err := db.producerToPush.SendMessage(ctx, key, &mqPushMsg)
return err
}
func (db *msgDatabase) MsgToMongoMQ(ctx context.Context, aggregationID string, triggerID string, messages []*pbMsg.MsgDataToMQ, lastSeq int64) error {
if len(messages) > 0 {
_, _, err := db.producerToModify.SendMessage(ctx, aggregationID, &pbMsg.MsgDataToMongoByMQ{LastSeq: lastSeq, AggregationID: aggregationID, Messages: messages, TriggerID: triggerID})
return err
}
return nil
}
2 years ago
func (db *msgDatabase) GetUserMaxSeq(ctx context.Context, userID string) (int64, error) {
2 years ago
return db.cache.GetUserMaxSeq(ctx, userID)
2 years ago
}
2 years ago
func (db *msgDatabase) GetUserMinSeq(ctx context.Context, userID string) (int64, error) {
2 years ago
return db.cache.GetUserMinSeq(ctx, userID)
2 years ago
}
2 years ago
func (db *msgDatabase) GetGroupMaxSeq(ctx context.Context, groupID string) (int64, error) {
2 years ago
return db.cache.GetGroupMaxSeq(ctx, groupID)
2 years ago
}
2 years ago
func (db *msgDatabase) GetGroupMinSeq(ctx context.Context, groupID string) (int64, error) {
2 years ago
return db.cache.GetGroupMinSeq(ctx, groupID)
2 years ago
}
2 years ago
func (db *msgDatabase) BatchInsertChat2DB(ctx context.Context, sourceID string, msgList []*pbMsg.MsgDataToMQ, currentMaxSeq int64) error {
2 years ago
//newTime := utils.GetCurrentTimestampByMill()
2 years ago
if int64(len(msgList)) > db.msg.GetSingleGocMsgNum() {
2 years ago
return errors.New("too large")
}
2 years ago
var remain int64
blk0 := db.msg.GetSingleGocMsgNum() - 1
2 years ago
//currentMaxSeq 4998
2 years ago
if currentMaxSeq < db.msg.GetSingleGocMsgNum() {
2 years ago
remain = blk0 - currentMaxSeq //1
} else {
excludeBlk0 := currentMaxSeq - blk0 //=1
//(5000-1)%5000 == 4999
2 years ago
remain = (db.msg.GetSingleGocMsgNum() - (excludeBlk0 % db.msg.GetSingleGocMsgNum())) % db.msg.GetSingleGocMsgNum()
2 years ago
}
//remain=1
2 years ago
var insertCounter int64
2 years ago
msgsToMongo := make([]unRelationTb.MsgInfoModel, 0)
msgsToMongoNext := make([]unRelationTb.MsgInfoModel, 0)
docID := ""
docIDNext := ""
var err error
for _, m := range msgList {
//log.Debug(operationID, "msg node ", m.String(), m.MsgData.ClientMsgID)
currentMaxSeq++
sMsg := unRelationTb.MsgInfoModel{}
sMsg.SendTime = m.MsgData.SendTime
2 years ago
m.MsgData.Seq = currentMaxSeq
2 years ago
if sMsg.Msg, err = proto.Marshal(m.MsgData); err != nil {
return utils.Wrap(err, "")
}
if insertCounter < remain {
msgsToMongo = append(msgsToMongo, sMsg)
insertCounter++
2 years ago
docID = db.msg.GetDocID(sourceID, currentMaxSeq)
2 years ago
//log.Debug(operationID, "msgListToMongo ", seqUid, m.MsgData.Seq, m.MsgData.ClientMsgID, insertCounter, remain, "userID: ", userID)
} else {
msgsToMongoNext = append(msgsToMongoNext, sMsg)
2 years ago
docIDNext = db.msg.GetDocID(sourceID, currentMaxSeq)
2 years ago
//log.Debug(operationID, "msgListToMongoNext ", seqUidNext, m.MsgData.Seq, m.MsgData.ClientMsgID, insertCounter, remain, "userID: ", userID)
}
}
if docID != "" {
//filter := bson.M{"uid": seqUid}
//log.NewDebug(operationID, "filter ", seqUid, "list ", msgListToMongo, "userID: ", userID)
//err := c.FindOneAndUpdate(ctx, filter, bson.M{"$push": bson.M{"msg": bson.M{"$each": msgsToMongo}}}).Err()
2 years ago
err = db.msgDocDatabase.PushMsgsToDoc(ctx, docID, msgsToMongo)
2 years ago
if err != nil {
if err == mongo.ErrNoDocuments {
doc := &unRelationTb.MsgDocModel{}
doc.DocID = docID
doc.Msg = msgsToMongo
2 years ago
if err = db.msgDocDatabase.Create(ctx, doc); err != nil {
2 years ago
prome.Inc(prome.MsgInsertMongoFailedCounter)
2 years ago
//log.NewError(operationID, "InsertOne failed", filter, err.Error(), sChat)
return utils.Wrap(err, "")
}
2 years ago
prome.Inc(prome.MsgInsertMongoSuccessCounter)
2 years ago
} else {
2 years ago
prome.Inc(prome.MsgInsertMongoFailedCounter)
2 years ago
//log.Error(operationID, "FindOneAndUpdate failed ", err.Error(), filter)
return utils.Wrap(err, "")
}
} else {
2 years ago
prome.Inc(prome.MsgInsertMongoSuccessCounter)
2 years ago
}
}
if docIDNext != "" {
nextDoc := &unRelationTb.MsgDocModel{}
nextDoc.DocID = docIDNext
nextDoc.Msg = msgsToMongoNext
//log.NewDebug(operationID, "filter ", seqUidNext, "list ", msgListToMongoNext, "userID: ", userID)
2 years ago
if err = db.msgDocDatabase.Create(ctx, nextDoc); err != nil {
2 years ago
prome.Inc(prome.MsgInsertMongoFailedCounter)
2 years ago
//log.NewError(operationID, "InsertOne failed", filter, err.Error(), sChat)
return utils.Wrap(err, "")
}
2 years ago
prome.Inc(prome.MsgInsertMongoSuccessCounter)
2 years ago
}
//log.Debug(operationID, "batch mgo cost time ", mongo2.getCurrentTimestampByMill()-newTime, userID, len(msgList))
return nil
}
2 years ago
func (db *msgDatabase) DeleteMessageFromCache(ctx context.Context, userID string, msgs []*pbMsg.MsgDataToMQ) error {
2 years ago
return db.cache.DeleteMessageFromCache(ctx, userID, msgs)
2 years ago
}
2 years ago
func (db *msgDatabase) BatchInsertChat2Cache(ctx context.Context, sourceID string, msgList []*pbMsg.MsgDataToMQ) (int64, error) {
2 years ago
//newTime := utils.GetCurrentTimestampByMill()
lenList := len(msgList)
2 years ago
if int64(lenList) > db.msg.GetSingleGocMsgNum() {
2 years ago
return 0, errors.New("too large")
}
if lenList < 1 {
return 0, errors.New("too short as 0")
}
// judge sessionType to get seq
2 years ago
var currentMaxSeq int64
2 years ago
var err error
if msgList[0].MsgData.SessionType == constant.SuperGroupChatType {
2 years ago
currentMaxSeq, err = db.cache.GetGroupMaxSeq(ctx, sourceID)
2 years ago
//log.Debug(operationID, "constant.SuperGroupChatType lastMaxSeq before add ", currentMaxSeq, "userID ", sourceID, err)
} else {
2 years ago
currentMaxSeq, err = db.cache.GetUserMaxSeq(ctx, sourceID)
2 years ago
//log.Debug(operationID, "constant.SingleChatType lastMaxSeq before add ", currentMaxSeq, "userID ", sourceID, err)
}
if err != nil && err != redis.Nil {
2 years ago
prome.Inc(prome.SeqGetFailedCounter)
2 years ago
return 0, utils.Wrap(err, "")
}
2 years ago
prome.Inc(prome.SeqGetSuccessCounter)
2 years ago
lastMaxSeq := currentMaxSeq
for _, m := range msgList {
currentMaxSeq++
2 years ago
m.MsgData.Seq = currentMaxSeq
2 years ago
//log.Debug(operationID, "cache msg node ", m.String(), m.MsgData.ClientMsgID, "userID: ", sourceID, "seq: ", currentMaxSeq)
}
//log.Debug(operationID, "SetMessageToCache ", sourceID, len(msgList))
2 years ago
failedNum, err := db.cache.SetMessageToCache(ctx, sourceID, msgList)
2 years ago
if err != nil {
2 years ago
prome.Add(prome.MsgInsertRedisFailedCounter, failedNum)
2 years ago
//log.Error(operationID, "setMessageToCache failed, continue ", err.Error(), len(msgList), sourceID)
} else {
2 years ago
prome.Inc(prome.MsgInsertRedisSuccessCounter)
2 years ago
}
//log.Debug(operationID, "batch to redis cost time ", mongo2.getCurrentTimestampByMill()-newTime, sourceID, len(msgList))
if msgList[0].MsgData.SessionType == constant.SuperGroupChatType {
2 years ago
err = db.cache.SetGroupMaxSeq(ctx, sourceID, currentMaxSeq)
2 years ago
} else {
2 years ago
err = db.cache.SetUserMaxSeq(ctx, sourceID, currentMaxSeq)
2 years ago
}
if err != nil {
2 years ago
prome.Inc(prome.SeqSetFailedCounter)
2 years ago
} else {
2 years ago
prome.Inc(prome.SeqSetSuccessCounter)
2 years ago
}
return lastMaxSeq, utils.Wrap(err, "")
}
2 years ago
func (db *msgDatabase) DelMsgBySeqs(ctx context.Context, userID string, seqs []int64) (totalUnExistSeqs []int64, err error) {
2 years ago
sortkeys.Int64s(seqs)
2 years ago
docIDSeqsMap := db.msg.GetDocIDSeqsMap(userID, seqs)
lock := sync.Mutex{}
var wg sync.WaitGroup
wg.Add(len(docIDSeqsMap))
for k, v := range docIDSeqsMap {
2 years ago
go func(docID string, seqs []int64) {
2 years ago
defer wg.Done()
unExistSeqList, err := db.DelMsgBySeqsInOneDoc(ctx, docID, seqs)
if err != nil {
return
}
lock.Lock()
totalUnExistSeqs = append(totalUnExistSeqs, unExistSeqList...)
lock.Unlock()
}(k, v)
}
return totalUnExistSeqs, nil
}
2 years ago
func (db *msgDatabase) DelMsgBySeqsInOneDoc(ctx context.Context, docID string, seqs []int64) (unExistSeqs []int64, err error) {
2 years ago
seqMsgs, indexes, unExistSeqs, err := db.GetMsgAndIndexBySeqsInOneDoc(ctx, docID, seqs)
if err != nil {
return nil, err
}
for i, v := range seqMsgs {
2 years ago
if err = db.msgDocDatabase.UpdateMsgStatusByIndexInOneDoc(ctx, docID, v, indexes[i], constant.MsgDeleted); err != nil {
2 years ago
return nil, err
}
}
return unExistSeqs, nil
}
2 years ago
func (db *msgDatabase) GetMsgAndIndexBySeqsInOneDoc(ctx context.Context, docID string, seqs []int64) (seqMsgs []*sdkws.MsgData, indexes []int, unExistSeqs []int64, err error) {
2 years ago
doc, err := db.msgDocDatabase.FindOneByDocID(ctx, docID)
2 years ago
if err != nil {
return nil, nil, nil, err
}
singleCount := 0
2 years ago
var hasSeqList []int64
2 years ago
for i := 0; i < len(doc.Msg); i++ {
msgPb, err := db.unmarshalMsg(&doc.Msg[i])
if err != nil {
return nil, nil, nil, err
}
2 years ago
if utils.Contain(msgPb.Seq, seqs...) {
2 years ago
indexes = append(indexes, i)
seqMsgs = append(seqMsgs, msgPb)
hasSeqList = append(hasSeqList, msgPb.Seq)
singleCount++
if singleCount == len(seqs) {
break
}
}
}
for _, i := range seqs {
2 years ago
if utils.Contain(i, hasSeqList...) {
2 years ago
continue
}
unExistSeqs = append(unExistSeqs, i)
}
return seqMsgs, indexes, unExistSeqs, nil
}
2 years ago
func (db *msgDatabase) GetNewestMsg(ctx context.Context, sourceID string) (msgPb *sdkws.MsgData, err error) {
2 years ago
msgInfo, err := db.msgDocDatabase.GetNewestMsg(ctx, sourceID)
2 years ago
if err != nil {
return nil, err
}
return db.unmarshalMsg(msgInfo)
}
2 years ago
func (db *msgDatabase) GetOldestMsg(ctx context.Context, sourceID string) (msgPb *sdkws.MsgData, err error) {
2 years ago
msgInfo, err := db.msgDocDatabase.GetOldestMsg(ctx, sourceID)
2 years ago
if err != nil {
return nil, err
}
return db.unmarshalMsg(msgInfo)
}
2 years ago
func (db *msgDatabase) unmarshalMsg(msgInfo *unRelationTb.MsgInfoModel) (msgPb *sdkws.MsgData, err error) {
2 years ago
msgPb = &sdkws.MsgData{}
err = proto.Unmarshal(msgInfo.Msg, msgPb)
if err != nil {
return nil, utils.Wrap(err, "")
}
return msgPb, nil
}
2 years ago
func (db *msgDatabase) getMsgBySeqs(ctx context.Context, sourceID string, seqs []int64, diffusionType int) (seqMsgs []*sdkws.MsgData, err error) {
2 years ago
var hasSeqs []int64
2 years ago
singleCount := 0
m := db.msg.GetDocIDSeqsMap(sourceID, seqs)
for docID, value := range m {
2 years ago
doc, err := db.msgDocDatabase.FindOneByDocID(ctx, docID)
2 years ago
if err != nil {
//log.NewError(operationID, "not find seqUid", seqUid, value, uid, seqList, err.Error())
continue
}
singleCount = 0
for i := 0; i < len(doc.Msg); i++ {
msgPb, err := db.unmarshalMsg(&doc.Msg[i])
if err != nil {
//log.NewError(operationID, "Unmarshal err", seqUid, value, uid, seqList, err.Error())
return nil, err
}
2 years ago
if utils.Contain(msgPb.Seq, value...) {
2 years ago
seqMsgs = append(seqMsgs, msgPb)
2 years ago
hasSeqs = append(hasSeqs, msgPb.Seq)
singleCount++
if singleCount == len(value) {
break
}
}
}
}
if len(hasSeqs) != len(seqs) {
2 years ago
var diff []int64
2 years ago
var exceptionMsg []*sdkws.MsgData
diff = utils.Difference(hasSeqs, seqs)
if diffusionType == constant.WriteDiffusion {
exceptionMsg = db.msg.GenExceptionMessageBySeqs(diff)
} else if diffusionType == constant.ReadDiffusion {
exceptionMsg = db.msg.GenExceptionSuperGroupMessageBySeqs(diff, sourceID)
}
2 years ago
seqMsgs = append(seqMsgs, exceptionMsg...)
2 years ago
}
2 years ago
return seqMsgs, nil
2 years ago
}
2 years ago
func (db *msgDatabase) GetMsgBySeqs(ctx context.Context, userID string, seqs []int64) (successMsgs []*sdkws.MsgData, err error) {
2 years ago
successMsgs, failedSeqs, err := db.cache.GetMessagesBySeq(ctx, userID, seqs)
2 years ago
if err != nil {
if err != redis.Nil {
2 years ago
prome.Add(prome.MsgPullFromRedisFailedCounter, len(failedSeqs))
2 years ago
log.Error(tracelog.GetOperationID(ctx), "get message from redis exception", err.Error(), failedSeqs)
}
}
2 years ago
prome.Add(prome.MsgPullFromRedisSuccessCounter, len(successMsgs))
2 years ago
if len(failedSeqs) > 0 {
mongoMsgs, err := db.getMsgBySeqs(ctx, userID, seqs, constant.WriteDiffusion)
if err != nil {
2 years ago
prome.Add(prome.MsgPullFromMongoFailedCounter, len(failedSeqs))
2 years ago
return nil, err
}
2 years ago
prome.Add(prome.MsgPullFromMongoSuccessCounter, len(mongoMsgs))
2 years ago
successMsgs = append(successMsgs, mongoMsgs...)
}
return successMsgs, nil
2 years ago
}
2 years ago
func (db *msgDatabase) GetSuperGroupMsgBySeqs(ctx context.Context, groupID string, seqs []int64) (successMsgs []*sdkws.MsgData, err error) {
2 years ago
successMsgs, failedSeqs, err := db.cache.GetMessagesBySeq(ctx, groupID, seqs)
2 years ago
if err != nil {
2 years ago
if err != redis.Nil {
2 years ago
prome.Add(prome.MsgPullFromRedisFailedCounter, len(failedSeqs))
2 years ago
log.Error(tracelog.GetOperationID(ctx), "get message from redis exception", err.Error(), failedSeqs)
}
2 years ago
}
2 years ago
prome.Add(prome.MsgPullFromRedisSuccessCounter, len(successMsgs))
2 years ago
if len(failedSeqs) > 0 {
mongoMsgs, err := db.getMsgBySeqs(ctx, groupID, seqs, constant.ReadDiffusion)
if err != nil {
2 years ago
prome.Add(prome.MsgPullFromMongoFailedCounter, len(failedSeqs))
2 years ago
return nil, err
}
2 years ago
prome.Add(prome.MsgPullFromMongoSuccessCounter, len(mongoMsgs))
2 years ago
successMsgs = append(successMsgs, mongoMsgs...)
2 years ago
}
2 years ago
return successMsgs, nil
}
2 years ago
func (db *msgDatabase) CleanUpUserMsg(ctx context.Context, userID string) error {
2 years ago
err := db.DeleteUserMsgsAndSetMinSeq(ctx, userID, 0)
2 years ago
if err != nil {
return err
}
2 years ago
err = db.cache.CleanUpOneUserAllMsg(ctx, userID)
2 years ago
return utils.Wrap(err, "")
2 years ago
}
2 years ago
func (db *msgDatabase) DeleteUserSuperGroupMsgsAndSetMinSeq(ctx context.Context, groupID string, userIDs []string, remainTime int64) error {
2 years ago
var delStruct delMsgRecursionStruct
minSeq, err := db.deleteMsgRecursion(ctx, groupID, unRelationTb.OldestList, &delStruct, remainTime)
if err != nil {
//log.NewError(operationID, utils.GetSelfFuncName(), groupID, "deleteMsg failed")
}
if minSeq == 0 {
return nil
}
//log.NewDebug(operationID, utils.GetSelfFuncName(), "delMsgIDList:", delStruct, "minSeq", minSeq)
for _, userID := range userIDs {
2 years ago
userMinSeq, err := db.cache.GetGroupUserMinSeq(ctx, groupID, userID)
2 years ago
if err != nil && err != redis.Nil {
//log.NewError(operationID, utils.GetSelfFuncName(), "GetGroupUserMinSeq failed", groupID, userID, err.Error())
continue
}
2 years ago
if userMinSeq > minSeq {
err = db.cache.SetGroupUserMinSeq(ctx, groupID, userID, userMinSeq)
2 years ago
} else {
2 years ago
err = db.cache.SetGroupUserMinSeq(ctx, groupID, userID, minSeq)
2 years ago
}
if err != nil {
//log.NewError(operationID, utils.GetSelfFuncName(), err.Error(), groupID, userID, userMinSeq, minSeq)
}
}
return nil
}
2 years ago
func (db *msgDatabase) DeleteUserMsgsAndSetMinSeq(ctx context.Context, userID string, remainTime int64) error {
2 years ago
var delStruct delMsgRecursionStruct
minSeq, err := db.deleteMsgRecursion(ctx, userID, unRelationTb.OldestList, &delStruct, remainTime)
if err != nil {
return utils.Wrap(err, "")
}
if minSeq == 0 {
return nil
}
2 years ago
return db.cache.SetUserMinSeq(ctx, userID, minSeq)
2 years ago
}
// this is struct for recursion
type delMsgRecursionStruct struct {
2 years ago
minSeq int64
delDocIDs []string
2 years ago
}
2 years ago
func (d *delMsgRecursionStruct) getSetMinSeq() int64 {
2 years ago
return d.minSeq
}
2 years ago
2 years ago
// index 0....19(del) 20...69
// seq 70
// set minSeq 21
// recursion 删除list并且返回设置的最小seq
2 years ago
func (db *msgDatabase) deleteMsgRecursion(ctx context.Context, sourceID string, index int64, delStruct *delMsgRecursionStruct, remainTime int64) (int64, error) {
2 years ago
// find from oldest list
2 years ago
msgs, err := db.msgDocDatabase.GetMsgsByIndex(ctx, sourceID, index)
2 years ago
if err != nil || msgs.DocID == "" {
if err != nil {
if err == unrelation.ErrMsgListNotExist {
2 years ago
log.NewDebug(tracelog.GetOperationID(ctx), utils.GetSelfFuncName(), "ID:", sourceID, "index:", index, err.Error())
2 years ago
} else {
//log.NewError(operationID, utils.GetSelfFuncName(), "GetUserMsgListByIndex failed", err.Error(), index, ID)
}
}
2 years ago
// 获取报错或者获取不到了物理删除并且返回seq delMongoMsgsPhysical(delStruct.delDocIDList), 结束递归
2 years ago
err = db.msgDocDatabase.Delete(ctx, delStruct.delDocIDs)
2 years ago
if err != nil {
return 0, err
}
return delStruct.getSetMinSeq() + 1, nil
}
//log.NewDebug(operationID, "ID:", sourceID, "index:", index, "uid:", msgs.UID, "len:", len(msgs.Msg))
2 years ago
if int64(len(msgs.Msg)) > db.msg.GetSingleGocMsgNum() {
2 years ago
log.NewWarn(tracelog.GetOperationID(ctx), utils.GetSelfFuncName(), "msgs too large:", len(msgs.Msg), "docID:", msgs.DocID)
}
if msgs.Msg[len(msgs.Msg)-1].SendTime+(remainTime*1000) < utils.GetCurrentTimestampByMill() && msgs.IsFull() {
2 years ago
delStruct.delDocIDs = append(delStruct.delDocIDs, msgs.DocID)
2 years ago
lastMsgPb := &sdkws.MsgData{}
err = proto.Unmarshal(msgs.Msg[len(msgs.Msg)-1].Msg, lastMsgPb)
if err != nil {
//log.NewError(operationID, utils.GetSelfFuncName(), err.Error(), len(msgs.Msg)-1, msgs.UID)
return 0, utils.Wrap(err, "proto.Unmarshal failed")
}
delStruct.minSeq = lastMsgPb.Seq
} else {
var hasMarkDelFlag bool
for _, msg := range msgs.Msg {
msgPb := &sdkws.MsgData{}
err = proto.Unmarshal(msg.Msg, msgPb)
if err != nil {
//log.NewError(operationID, utils.GetSelfFuncName(), err.Error(), len(msgs.Msg)-1, msgs.UID)
return 0, utils.Wrap(err, "proto.Unmarshal failed")
}
if utils.GetCurrentTimestampByMill() > msg.SendTime+(remainTime*1000) {
msgPb.Status = constant.MsgDeleted
bytes, _ := proto.Marshal(msgPb)
msg.Msg = bytes
msg.SendTime = 0
hasMarkDelFlag = true
} else {
2 years ago
// 到本条消息不需要删除, minSeq置为这条消息的seq
2 years ago
if err := db.msgDocDatabase.Delete(ctx, delStruct.delDocIDs); err != nil {
2 years ago
return 0, err
}
if hasMarkDelFlag {
2 years ago
if err := db.msgDocDatabase.UpdateOneDoc(ctx, msgs); err != nil {
2 years ago
return delStruct.getSetMinSeq(), utils.Wrap(err, "")
}
}
2 years ago
return msgPb.Seq, nil
2 years ago
}
}
}
//log.NewDebug(operationID, sourceID, "continue to", delStruct)
// 继续递归 index+1
seq, err := db.deleteMsgRecursion(ctx, sourceID, index+1, delStruct, remainTime)
return seq, utils.Wrap(err, "deleteMsg failed")
2 years ago
}
2 years ago
2 years ago
func (db *msgDatabase) GetUserMinMaxSeqInMongoAndCache(ctx context.Context, userID string) (minSeqMongo, maxSeqMongo, minSeqCache, maxSeqCache int64, err error) {
2 years ago
minSeqMongo, maxSeqMongo, err = db.GetMinMaxSeqMongo(ctx, userID)
if err != nil {
return 0, 0, 0, 0, err
}
// from cache
minSeqCache, err = db.cache.GetUserMinSeq(ctx, userID)
if err != nil {
return 0, 0, 0, 0, err
}
maxSeqCache, err = db.cache.GetUserMaxSeq(ctx, userID)
if err != nil {
return 0, 0, 0, 0, err
}
return
}
2 years ago
func (db *msgDatabase) GetSuperGroupMinMaxSeqInMongoAndCache(ctx context.Context, groupID string) (minSeqMongo, maxSeqMongo, maxSeqCache int64, err error) {
2 years ago
minSeqMongo, maxSeqMongo, err = db.GetMinMaxSeqMongo(ctx, groupID)
if err != nil {
return 0, 0, 0, err
}
maxSeqCache, err = db.cache.GetGroupMaxSeq(ctx, groupID)
if err != nil {
return 0, 0, 0, err
}
return
}
2 years ago
func (db *msgDatabase) GetMinMaxSeqMongo(ctx context.Context, sourceID string) (minSeqMongo, maxSeqMongo int64, err error) {
2 years ago
oldestMsgMongo, err := db.msgDocDatabase.GetOldestMsg(ctx, sourceID)
2 years ago
if err != nil {
return 0, 0, err
}
msgPb, err := db.unmarshalMsg(oldestMsgMongo)
if err != nil {
return 0, 0, err
}
minSeqMongo = msgPb.Seq
2 years ago
newestMsgMongo, err := db.msgDocDatabase.GetNewestMsg(ctx, sourceID)
2 years ago
if err != nil {
return 0, 0, err
}
msgPb, err = db.unmarshalMsg(newestMsgMongo)
if err != nil {
return 0, 0, err
}
maxSeqMongo = msgPb.Seq
return
}
2 years ago
func (db *msgDatabase) SetGroupUserMinSeq(ctx context.Context, groupID, userID string, minSeq int64) (err error) {
2 years ago
return db.cache.SetGroupUserMinSeq(ctx, groupID, userID, minSeq)
}
2 years ago
func (db *msgDatabase) SetUserMinSeq(ctx context.Context, userID string, minSeq int64) (err error) {
2 years ago
return db.cache.SetUserMinSeq(ctx, userID, minSeq)
}
2 years ago
func (db *msgDatabase) GetGroupUserMinSeq(ctx context.Context, groupID, userID string) (int64, error) {
return db.cache.GetGroupUserMinSeq(ctx, groupID, userID)
}