@@ -11,7 +11,6 @@ import (
"github.com/gogo/protobuf/sortkeys"
"sync"
//"Open_IM/pkg/common/log"
pbMsg "Open_IM/pkg/proto/msg"
"Open_IM/pkg/proto/sdkws"
"Open_IM/pkg/utils"
@@ -25,30 +24,31 @@ import (
type MsgInterface interface {
// 批量插入消息到db
BatchInsertChat2DB ( ctx context . Context , ID string , msgList [ ] * pbMsg . MsgDataToMQ , currentMaxSeq u int64) error
BatchInsertChat2DB ( ctx context . Context , source ID string , msgList [ ] * pbMsg . MsgDataToMQ , currentMaxSeq int64 ) error
// 刪除redis中消息缓存
DeleteMessageFromCache ( ctx context . Context , user ID string , msgList [ ] * pbMsg . MsgDataToMQ ) error
DeleteMessageFromCache ( ctx context . Context , source ID string , msgList [ ] * pbMsg . MsgDataToMQ ) error
// incrSeq然后批量插入缓存
BatchInsertChat2Cache ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ ) ( u int64, error )
BatchInsertChat2Cache ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ ) ( int64 , error )
// 删除消息 返回不存在的seqList
DelMsgBySeqs ( ctx context . Context , userID string , seqs [ ] u int32 ) ( totalUnExistSeqs [ ] u int32 , err error )
// 获取群ID或者UserID最新一条在db里面的消息
GetNewestMsg ( ctx context . Context , sourceID string ) ( msg * sdkws . MsgData , err error )
// 获取群ID或者UserID最老一条在db里面的消息
GetOldestMsg ( ctx context . Context , sourceID string ) ( msg * sdkws . MsgData , err error )
DelMsgBySeqs ( ctx context . Context , userID string , seqs [ ] int64 ) ( totalUnExistSeqs [ ] int64 , err error )
// 通过seqList获取db中写扩散消息
GetMsgBySeqs ( ctx context . Context , userID string , seqs [ ] u int32 ) ( seqMsg [ ] * sdkws . MsgData , err error )
GetMsgBySeqs ( ctx context . Context , userID string , seqs [ ] int64 ) ( seqMsg [ ] * sdkws . MsgData , err error )
// 通过seqList获取大群在db里面的消息
GetSuperGroupMsgBySeqs ( ctx context . Context , groupID string , seqs [ ] u int32 ) ( seqMsg [ ] * sdkws . MsgData , err error )
GetSuperGroupMsgBySeqs ( ctx context . Context , groupID string , seqs [ ] int64 ) ( seqMsg [ ] * sdkws . MsgData , err error )
// 删除用户所有消息/cache/db然后重置seq
CleanUpUserMsg ( ctx context . Context , userID string ) error
// 删除大群消息重置群成员最小群seq, remainTime为消息保留的时间单位秒,超时消息删除, 传0删除所有消息(此方法不删除 redis cache)
DeleteUserSuperGroupMsgsAndSetMinSeq ( ctx context . Context , groupID string , userID string , remainTime int64 ) error
DeleteUserSuperGroupMsgsAndSetMinSeq ( ctx context . Context , groupID string , userID [ ] string , remainTime int64 ) error
// 删除用户消息重置最小seq, remainTime为消息保留的时间单位秒,超时消息删除, 传0删除所有消息(此方法不删除redis cache)
DeleteUserMsgsAndSetMinSeq ( ctx context . Context , userID string , remainTime int64 ) error
// SetSendMsgStatus
// GetSendMsgStatu s
// 获取用户 seq mongo和redis
GetUserMinMaxSeqInMongoAndCache ( ctx context . Context , userID string ) ( minSeqMongo , maxSeqMongo , minSeqCache , maxSeqCache int64 , err error )
// 获取群 seq mongo和redi s
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 )
// 设置用户最小seq 直接调用cache
SetUserMinSeq ( ctx context . Context , userID string , minSeq int64 ) ( err error )
}
func NewMsgController ( mgo * mongo . Client , rdb redis . UniversalClient ) MsgInterface {
@@ -59,35 +59,27 @@ type MsgController struct {
database MsgDatabase
}
func ( m * MsgController ) BatchInsertChat2DB ( ctx context . Context , ID string , msgList [ ] * pbMsg . MsgDataToMQ , currentMaxSeq u int64) error {
func ( m * MsgController ) BatchInsertChat2DB ( ctx context . Context , ID string , msgList [ ] * pbMsg . MsgDataToMQ , currentMaxSeq int64 ) error {
return m . database . BatchInsertChat2DB ( ctx , ID , msgList , currentMaxSeq )
}
func ( m * MsgController ) DeleteMessageFromCache ( ctx context . Context , user ID string , msgList [ ] * pbMsg . MsgDataToMQ ) error {
return m . database . DeleteMessageFromCache ( ctx , user ID, msgList )
func ( m * MsgController ) DeleteMessageFromCache ( ctx context . Context , source ID string , msgList [ ] * pbMsg . MsgDataToMQ ) error {
return m . database . DeleteMessageFromCache ( ctx , source ID, msgList )
}
func ( m * MsgController ) BatchInsertChat2Cache ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ ) ( u int64, error ) {
func ( m * MsgController ) BatchInsertChat2Cache ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ ) ( int64 , error ) {
return m . database . BatchInsertChat2Cache ( ctx , sourceID , msgList )
}
func ( m * MsgController ) DelMsgBySeqs ( ctx context . Context , userID string , seqs [ ] u int32 ) ( totalUnExistSeqs [ ] u int32 , err error ) {
func ( m * MsgController ) DelMsgBySeqs ( ctx context . Context , userID string , seqs [ ] int64 ) ( totalUnExistSeqs [ ] int64 , err error ) {
return m . database . DelMsgBySeqs ( ctx , userID , seqs )
}
func ( m * MsgController ) GetNewestMsg ( ctx context . Context , ID string ) ( msg * sdkws . MsgData , err error ) {
return m . database . GetNewestMsg ( ctx , ID )
}
func ( m * MsgController ) GetOldestMsg ( ctx context . Context , ID string ) ( msg * sdkws . MsgData , err error ) {
return m . database . GetOldestMsg ( ctx , ID )
}
func ( m * MsgController ) GetMsgBySeqs ( ctx context . Context , userID string , seqs [ ] uint32 ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
func ( m * MsgController ) GetMsgBySeqs ( ctx context . Context , user ID string , seqs [ ] int64 ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
return m . database . GetMsgBySeqs ( ctx , userID , seqs )
}
func ( m * MsgController ) GetSuperGroupMsgBySeqs ( ctx context . Context , groupID string , seqs [ ] u int32 ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
func ( m * MsgController ) GetSuperGroupMsgBySeqs ( ctx context . Context , groupID string , seqs [ ] int64 ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
return m . database . GetSuperGroupMsgBySeqs ( ctx , groupID , seqs )
}
@@ -95,66 +87,87 @@ func (m *MsgController) CleanUpUserMsg(ctx context.Context, userID string) error
return m . database . CleanUpUserMsg ( ctx , userID )
}
func ( m * MsgController ) DeleteUserSuperGroupMsgsAndSetMinSeq ( ctx context . Context , groupID string , userID string , remainTime int64 ) error {
return m . database . DeleteUserMsgsAndSetMinSeq ( ctx , userID , remainTime )
func ( m * MsgController ) DeleteUserSuperGroupMsgsAndSetMinSeq ( ctx context . Context , groupID string , userIDs [ ] string , remainTime int64 ) error {
return m . database . DeleteUserSuperGroup MsgsAndSetMinSeq ( ctx , groupID , userIDs , remainTime )
}
func ( m * MsgController ) DeleteUserMsgsAndSetMinSeq ( ctx context . Context , userID string , remainTime int64 ) error {
return m . database . DeleteUserMsgsAndSetMinSeq ( ctx , userID , remainTime )
}
func ( m * MsgController ) GetUserMinMaxSeqInMongoAndCache ( ctx context . Context , userID string ) ( minSeqMongo , maxSeqMongo , minSeqCache , maxSeqCache int64 , err error ) {
return m . database . GetUserMinMaxSeqInMongoAndCache ( ctx , userID )
}
func ( m * MsgController ) GetSuperGroupMinMaxSeqInMongoAndCache ( ctx context . Context , groupID string ) ( minSeqMongo , maxSeqMongo , maxSeqCache int64 , err error ) {
return m . database . GetSuperGroupMinMaxSeqInMongoAndCache ( ctx , groupID )
}
func ( m * MsgController ) SetGroupUserMinSeq ( ctx context . Context , groupID , userID string , minSeq int64 ) ( err error ) {
return m . database . SetGroupUserMinSeq ( ctx , groupID , userID , minSeq )
}
func ( m * MsgController ) SetUserMinSeq ( ctx context . Context , userID string , minSeq int64 ) ( err error ) {
return m . database . SetUserMinSeq ( ctx , userID , minSeq )
}
type MsgDatabaseInterface interface {
// 批量插入消息
BatchInsertChat2DB ( ctx context . Context , ID string , msgList [ ] * pbMsg . MsgDataToMQ , currentMaxSeq u int64) error
BatchInsertChat2DB ( ctx context . Context , source ID string , msgList [ ] * pbMsg . MsgDataToMQ , currentMaxSeq int64 ) error
// 刪除redis中消息缓存
DeleteMessageFromCache ( ctx context . Context , user ID string , msgList [ ] * pbMsg . MsgDataToMQ ) error
DeleteMessageFromCache ( ctx context . Context , source ID string , msgList [ ] * pbMsg . MsgDataToMQ ) error
// incrSeq然后批量插入缓存
BatchInsertChat2Cache ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ ) ( u int64, error )
BatchInsertChat2Cache ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ ) ( int64 , error )
// 删除消息 返回不存在的seqList
DelMsgBySeqs ( ctx context . Context , userID string , seqs [ ] u int32 ) ( totalUnExistSeqs [ ] u int32 , err error )
DelMsgBySeqs ( ctx context . Context , userID string , seqs [ ] int64 ) ( totalUnExistSeqs [ ] int64 , err error )
// 获取群ID或者UserID最新一条在mongo里面的消息
GetNewestMsg ( ctx context . Context , sourceID string ) ( msg * sdkws . MsgData , err error )
// 获取群ID或者UserID最老一条在mongo里面的消息
GetOldestMsg ( ctx context . Context , sourceID string ) ( msg * sdkws . MsgData , err error )
// 通过seqList获取mongo中写扩散消息
GetMsgBySeqs ( ctx context . Context , userID string , seqs [ ] u int32 ) ( seqMsg [ ] * sdkws . MsgData , err error )
GetMsgBySeqs ( ctx context . Context , userID string , seqs [ ] int64 ) ( seqMsg [ ] * sdkws . MsgData , err error )
// 通过seqList获取大群在 mongo里面的消息
GetSuperGroupMsgBySeqs ( ctx context . Context , groupID string , seqs [ ] u int32 ) ( seqMsg [ ] * sdkws . MsgData , err error )
GetSuperGroupMsgBySeqs ( ctx context . Context , groupID string , seqs [ ] int64 ) ( seqMsg [ ] * sdkws . MsgData , err error )
// 删除用户所有消息/redis/mongo然后重置seq
CleanUpUserMsg ( ctx context . Context , userID string ) error
// 删除大群消息重置群成员最小群seq, remainTime为消息保留的时间单位秒,超时消息删除, 传0删除所有消息(此方法不删除 redis cache)
DeleteUserSuperGroupMsgsAndSetMinSeq ( ctx context . Context , groupID string , userID [ ] string , remainTime int64 ) error
DeleteUserSuperGroupMsgsAndSetMinSeq ( ctx context . Context , groupID string , userIDs [ ] string , remainTime int64 ) error
// 删除用户消息重置最小seq, remainTime为消息保留的时间单位秒,超时消息删除, 传0删除所有消息(此方法不删除redis cache)
DeleteUserMsgsAndSetMinSeq ( ctx context . Context , userID string , remainTime int64 ) error
// 获取用户 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 )
// 设置用户最小seq 直接调用cache
SetUserMinSeq ( ctx context . Context , userID string , minSeq int64 ) ( err error )
}
type MsgDatabase struct {
msgModel unRelationTb . MsgDocModelInterface
msgC ache cache . Cache
msg unRelationTb . MsgDocModel
mgo unRelationTb . MsgDocModelInterface
c ache cache . Cache
msg unRelationTb . MsgDocModel
}
func NewMsgDatabase ( mgo * mongo . Client , rdb redis . UniversalClient ) MsgDatabaseInterface {
return & MsgDatabase { }
}
func ( db * MsgDatabase ) BatchInsertChat2DB ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ , currentMaxSeq u int64) error {
func ( db * MsgDatabase ) BatchInsertChat2DB ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ , currentMaxSeq int64 ) error {
//newTime := utils.GetCurrentTimestampByMill()
if len ( msgList ) > db . msg . GetSingleGocMsgNum ( ) {
if int64 ( len( msgList ) ) > db . msg . GetSingleGocMsgNum ( ) {
return errors . New ( "too large" )
}
var remain u int64
blk0 := uint64 ( db . msg . GetSingleGocMsgNum ( ) - 1 )
var remain int64
blk0 := db . msg . GetSingleGocMsgNum ( ) - 1
//currentMaxSeq 4998
if currentMaxSeq < uint64 ( db . msg . GetSingleGocMsgNum ( ) ) {
if currentMaxSeq < db . msg . GetSingleGocMsgNum ( ) {
remain = blk0 - currentMaxSeq //1
} else {
excludeBlk0 := currentMaxSeq - blk0 //=1
//(5000-1)%5000 == 4999
remain = ( uint64 ( db . msg . GetSingleGocMsgNum ( ) ) - ( excludeBlk0 % uint64 ( db . msg . GetSingleGocMsgNum ( ) ) ) ) % uint64 ( db . msg . GetSingleGocMsgNum ( ) )
remain = ( db . msg . GetSingleGocMsgNum ( ) - ( excludeBlk0 % db . msg . GetSingleGocMsgNum ( ) ) ) % db . msg . GetSingleGocMsgNum ( )
}
//remain=1
insertCounter := uint64 ( 0 )
var insertCounter int64
msgsToMongo := make ( [ ] unRelationTb . MsgInfoModel , 0 )
msgsToMongoNext := make ( [ ] unRelationTb . MsgInfoModel , 0 )
docID := ""
@@ -165,18 +178,18 @@ func (db *MsgDatabase) BatchInsertChat2DB(ctx context.Context, sourceID string,
currentMaxSeq ++
sMsg := unRelationTb . MsgInfoModel { }
sMsg . SendTime = m . MsgData . SendTime
m . MsgData . Seq = uint32 ( currentMaxSeq )
m . MsgData . Seq = currentMaxSeq
if sMsg . Msg , err = proto . Marshal ( m . MsgData ) ; err != nil {
return utils . Wrap ( err , "" )
}
if insertCounter < remain {
msgsToMongo = append ( msgsToMongo , sMsg )
insertCounter ++
docID = db . msg . GetDocID ( sourceID , uint32 ( currentMaxSeq ) )
docID = db . msg . GetDocID ( sourceID , currentMaxSeq )
//log.Debug(operationID, "msgListToMongo ", seqUid, m.MsgData.Seq, m.MsgData.ClientMsgID, insertCounter, remain, "userID: ", userID)
} else {
msgsToMongoNext = append ( msgsToMongoNext , sMsg )
docIDNext = db . msg . GetDocID ( sourceID , uint32 ( currentMaxSeq ) )
docIDNext = db . msg . GetDocID ( sourceID , currentMaxSeq )
//log.Debug(operationID, "msgListToMongoNext ", seqUidNext, m.MsgData.Seq, m.MsgData.ClientMsgID, insertCounter, remain, "userID: ", userID)
}
}
@@ -185,13 +198,13 @@ func (db *MsgDatabase) BatchInsertChat2DB(ctx context.Context, sourceID string,
//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()
err = db . msgModel . PushMsgsToDoc ( ctx , docID , msgsToMongo )
err = db . mgo . PushMsgsToDoc ( ctx , docID , msgsToMongo )
if err != nil {
if err == mongo . ErrNoDocuments {
doc := & unRelationTb . MsgDocModel { }
doc . DocID = docID
doc . Msg = msgsToMongo
if err = db . msgModel . Create ( ctx , doc ) ; err != nil {
if err = db . mgo . Create ( ctx , doc ) ; err != nil {
prome . PromeInc ( prome . MsgInsertMongoFailedCounter )
//log.NewError(operationID, "InsertOne failed", filter, err.Error(), sChat)
return utils . Wrap ( err , "" )
@@ -211,7 +224,7 @@ func (db *MsgDatabase) BatchInsertChat2DB(ctx context.Context, sourceID string,
nextDoc . DocID = docIDNext
nextDoc . Msg = msgsToMongoNext
//log.NewDebug(operationID, "filter ", seqUidNext, "list ", msgListToMongoNext, "userID: ", userID)
if err = db . msgModel . Create ( ctx , nextDoc ) ; err != nil {
if err = db . mgo . Create ( ctx , nextDoc ) ; err != nil {
prome . PromeInc ( prome . MsgInsertMongoFailedCounter )
//log.NewError(operationID, "InsertOne failed", filter, err.Error(), sChat)
return utils . Wrap ( err , "" )
@@ -223,26 +236,26 @@ func (db *MsgDatabase) BatchInsertChat2DB(ctx context.Context, sourceID string,
}
func ( db * MsgDatabase ) DeleteMessageFromCache ( ctx context . Context , userID string , msgs [ ] * pbMsg . MsgDataToMQ ) error {
return db . msgC ache. DeleteMessageFromCache ( ctx , userID , msgs )
return db . c ache. DeleteMessageFromCache ( ctx , userID , msgs )
}
func ( db * MsgDatabase ) BatchInsertChat2Cache ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ ) ( u int64, error ) {
func ( db * MsgDatabase ) BatchInsertChat2Cache ( ctx context . Context , sourceID string , msgList [ ] * pbMsg . MsgDataToMQ ) ( int64 , error ) {
//newTime := utils.GetCurrentTimestampByMill()
lenList := len ( msgList )
if lenList > db . msg . GetSingleGocMsgNum ( ) {
if int64 ( lenList ) > db . msg . GetSingleGocMsgNum ( ) {
return 0 , errors . New ( "too large" )
}
if lenList < 1 {
return 0 , errors . New ( "too short as 0" )
}
// judge sessionType to get seq
var currentMaxSeq u int64
var currentMaxSeq int64
var err error
if msgList [ 0 ] . MsgData . SessionType == constant . SuperGroupChatType {
currentMaxSeq , err = db . msgC ache. GetGroupMaxSeq ( ctx , sourceID )
currentMaxSeq , err = db . c ache. GetGroupMaxSeq ( ctx , sourceID )
//log.Debug(operationID, "constant.SuperGroupChatType lastMaxSeq before add ", currentMaxSeq, "userID ", sourceID, err)
} else {
currentMaxSeq , err = db . msgC ache. GetUserMaxSeq ( ctx , sourceID )
currentMaxSeq , err = db . c ache. GetUserMaxSeq ( ctx , sourceID )
//log.Debug(operationID, "constant.SingleChatType lastMaxSeq before add ", currentMaxSeq, "userID ", sourceID, err)
}
if err != nil && err != redis . Nil {
@@ -253,11 +266,11 @@ func (db *MsgDatabase) BatchInsertChat2Cache(ctx context.Context, sourceID strin
lastMaxSeq := currentMaxSeq
for _ , m := range msgList {
currentMaxSeq ++
m . MsgData . Seq = uint32 ( currentMaxSeq )
m . MsgData . Seq = currentMaxSeq
//log.Debug(operationID, "cache msg node ", m.String(), m.MsgData.ClientMsgID, "userID: ", sourceID, "seq: ", currentMaxSeq)
}
//log.Debug(operationID, "SetMessageToCache ", sourceID, len(msgList))
failedNum , err := db . msgC ache. SetMessageToCache ( ctx , sourceID , msgList )
failedNum , err := db . c ache. SetMessageToCache ( ctx , sourceID , msgList )
if err != nil {
prome . PromeAdd ( prome . MsgInsertRedisFailedCounter , failedNum )
//log.Error(operationID, "setMessageToCache failed, continue ", err.Error(), len(msgList), sourceID)
@@ -266,9 +279,9 @@ func (db *MsgDatabase) BatchInsertChat2Cache(ctx context.Context, sourceID strin
}
//log.Debug(operationID, "batch to redis cost time ", mongo2.getCurrentTimestampByMill()-newTime, sourceID, len(msgList))
if msgList [ 0 ] . MsgData . SessionType == constant . SuperGroupChatType {
err = db . msgC ache. SetGroupMaxSeq ( ctx , sourceID , currentMaxSeq )
err = db . c ache. SetGroupMaxSeq ( ctx , sourceID , currentMaxSeq )
} else {
err = db . msgC ache. SetUserMaxSeq ( ctx , sourceID , currentMaxSeq )
err = db . c ache. SetUserMaxSeq ( ctx , sourceID , currentMaxSeq )
}
if err != nil {
prome . PromeInc ( prome . SeqSetFailedCounter )
@@ -278,14 +291,14 @@ func (db *MsgDatabase) BatchInsertChat2Cache(ctx context.Context, sourceID strin
return lastMaxSeq , utils . Wrap ( err , "" )
}
func ( db * MsgDatabase ) DelMsgBySeqs ( ctx context . Context , userID string , seqs [ ] u int32 ) ( totalUnExistSeqs [ ] u int32 , err error ) {
sortkeys . Uint32 s( seqs )
func ( db * MsgDatabase ) DelMsgBySeqs ( ctx context . Context , userID string , seqs [ ] int64 ) ( totalUnExistSeqs [ ] int64 , err error ) {
sortkeys . Int64 s( seqs )
docIDSeqsMap := db . msg . GetDocIDSeqsMap ( userID , seqs )
lock := sync . Mutex { }
var wg sync . WaitGroup
wg . Add ( len ( docIDSeqsMap ) )
for k , v := range docIDSeqsMap {
go func ( docID string , seqs [ ] u int32 ) {
go func ( docID string , seqs [ ] int64 ) {
defer wg . Done ( )
unExistSeqList , err := db . DelMsgBySeqsInOneDoc ( ctx , docID , seqs )
if err != nil {
@@ -299,26 +312,26 @@ func (db *MsgDatabase) DelMsgBySeqs(ctx context.Context, userID string, seqs []u
return totalUnExistSeqs , nil
}
func ( db * MsgDatabase ) DelMsgBySeqsInOneDoc ( ctx context . Context , docID string , seqs [ ] u int32 ) ( unExistSeqs [ ] u int32 , err error ) {
func ( db * MsgDatabase ) DelMsgBySeqsInOneDoc ( ctx context . Context , docID string , seqs [ ] int64 ) ( unExistSeqs [ ] int64 , err error ) {
seqMsgs , indexes , unExistSeqs , err := db . GetMsgAndIndexBySeqsInOneDoc ( ctx , docID , seqs )
if err != nil {
return nil , err
}
for i , v := range seqMsgs {
if err = db . msgModel . UpdateMsgStatusByIndexInOneDoc ( ctx , docID , v , indexes [ i ] , constant . MsgDeleted ) ; err != nil {
if err = db . mgo . UpdateMsgStatusByIndexInOneDoc ( ctx , docID , v , indexes [ i ] , constant . MsgDeleted ) ; err != nil {
return nil , err
}
}
return unExistSeqs , nil
}
func ( db * MsgDatabase ) GetMsgAndIndexBySeqsInOneDoc ( ctx context . Context , docID string , seqs [ ] u int32 ) ( seqMsgs [ ] * sdkws . MsgData , indexes [ ] int , unExistSeqs [ ] u int32 , err error ) {
doc , err := db . msgModel . FindOneByDocID ( ctx , docID )
func ( db * MsgDatabase ) GetMsgAndIndexBySeqsInOneDoc ( ctx context . Context , docID string , seqs [ ] int64 ) ( seqMsgs [ ] * sdkws . MsgData , indexes [ ] int , unExistSeqs [ ] int64 , err error ) {
doc , err := db . mgo . FindOneByDocID ( ctx , docID )
if err != nil {
return nil , nil , nil , err
}
singleCount := 0
var hasSeqList [ ] u int32
var hasSeqList [ ] int64
for i := 0 ; i < len ( doc . Msg ) ; i ++ {
msgPb , err := db . unmarshalMsg ( & doc . Msg [ i ] )
if err != nil {
@@ -344,7 +357,7 @@ func (db *MsgDatabase) GetMsgAndIndexBySeqsInOneDoc(ctx context.Context, docID s
}
func ( db * MsgDatabase ) GetNewestMsg ( ctx context . Context , sourceID string ) ( msgPb * sdkws . MsgData , err error ) {
msgInfo , err := db . msgModel . GetNewestMsg ( ctx , sourceID )
msgInfo , err := db . mgo . GetNewestMsg ( ctx , sourceID )
if err != nil {
return nil , err
}
@@ -352,7 +365,7 @@ func (db *MsgDatabase) GetNewestMsg(ctx context.Context, sourceID string) (msgPb
}
func ( db * MsgDatabase ) GetOldestMsg ( ctx context . Context , sourceID string ) ( msgPb * sdkws . MsgData , err error ) {
msgInfo , err := db . msgModel . GetOldestMsg ( ctx , sourceID )
msgInfo , err := db . mgo . GetOldestMsg ( ctx , sourceID )
if err != nil {
return nil , err
}
@@ -368,12 +381,12 @@ func (db *MsgDatabase) unmarshalMsg(msgInfo *unRelationTb.MsgInfoModel) (msgPb *
return msgPb , nil
}
func ( db * MsgDatabase ) getMsgBySeqs ( ctx context . Context , sourceID string , seqs [ ] u int32 , diffusionType int ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
var hasSeqs [ ] u int32
func ( db * MsgDatabase ) getMsgBySeqs ( ctx context . Context , sourceID string , seqs [ ] int64 , diffusionType int ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
var hasSeqs [ ] int64
singleCount := 0
m := db . msg . GetDocIDSeqsMap ( sourceID , seqs )
for docID , value := range m {
doc , err := db . msgModel . FindOneByDocID ( ctx , docID )
doc , err := db . mgo . FindOneByDocID ( ctx , docID )
if err != nil {
//log.NewError(operationID, "not find seqUid", seqUid, value, uid, seqList, err.Error())
continue
@@ -396,7 +409,7 @@ func (db *MsgDatabase) getMsgBySeqs(ctx context.Context, sourceID string, seqs [
}
}
if len ( hasSeqs ) != len ( seqs ) {
var diff [ ] u int32
var diff [ ] int64
var exceptionMsg [ ] * sdkws . MsgData
diff = utils . Difference ( hasSeqs , seqs )
if diffusionType == constant . WriteDiffusion {
@@ -409,8 +422,8 @@ func (db *MsgDatabase) getMsgBySeqs(ctx context.Context, sourceID string, seqs [
return seqMsg , nil
}
func ( db * MsgDatabase ) GetMsgBySeqs ( ctx context . Context , userID string , seqs [ ] u int32 ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
successMsgs , failedSeqs , err := db . msgC ache. GetMessageListBySeq ( ctx , userID , seqs )
func ( db * MsgDatabase ) GetMsgBySeqs ( ctx context . Context , userID string , seqs [ ] int64 ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
successMsgs , failedSeqs , err := db . c ache. GetMessageListBySeq ( ctx , userID , seqs )
if err != nil {
if err != redis . Nil {
prome . PromeAdd ( prome . MsgPullFromRedisFailedCounter , len ( failedSeqs ) )
@@ -430,8 +443,8 @@ func (db *MsgDatabase) GetMsgBySeqs(ctx context.Context, userID string, seqs []u
return successMsgs , nil
}
func ( db * MsgDatabase ) GetSuperGroupMsgBySeqs ( ctx context . Context , groupID string , seqs [ ] u int32 ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
successMsgs , failedSeqs , err := db . msgC ache. GetMessageListBySeq ( ctx , groupID , seqs )
func ( db * MsgDatabase ) GetSuperGroupMsgBySeqs ( ctx context . Context , groupID string , seqs [ ] int64 ) ( seqMsg [ ] * sdkws . MsgData , err error ) {
successMsgs , failedSeqs , err := db . c ache. GetMessageListBySeq ( ctx , groupID , seqs )
if err != nil {
if err != redis . Nil {
prome . PromeAdd ( prome . MsgPullFromRedisFailedCounter , len ( failedSeqs ) )
@@ -456,7 +469,7 @@ func (db *MsgDatabase) CleanUpUserMsg(ctx context.Context, userID string) error
if err != nil {
return err
}
err = db . msgC ache. CleanUpOneUserAllMsg ( ctx , userID )
err = db . c ache. CleanUpOneUserAllMsg ( ctx , userID )
return utils . Wrap ( err , "" )
}
@@ -471,15 +484,15 @@ func (db *MsgDatabase) DeleteUserSuperGroupMsgsAndSetMinSeq(ctx context.Context,
}
//log.NewDebug(operationID, utils.GetSelfFuncName(), "delMsgIDList:", delStruct, "minSeq", minSeq)
for _ , userID := range userIDs {
userMinSeq , err := db . msgC ache. GetGroupUserMinSeq ( ctx , groupID , userID )
userMinSeq , err := db . c ache. GetGroupUserMinSeq ( ctx , groupID , userID )
if err != nil && err != redis . Nil {
//log.NewError(operationID, utils.GetSelfFuncName(), "GetGroupUserMinSeq failed", groupID, userID, err.Error())
continue
}
if userMinSeq > uint64 ( minSeq ) {
err = db . msgC ache. SetGroupUserMinSeq ( ctx , groupID , userID , userMinSeq )
if userMinSeq > minSeq {
err = db . c ache. SetGroupUserMinSeq ( ctx , groupID , userID , userMinSeq )
} else {
err = db . msgC ache. SetGroupUserMinSeq ( ctx , groupID , userID , uint64 ( minSeq ) )
err = db . c ache. SetGroupUserMinSeq ( ctx , groupID , userID , minSeq )
}
if err != nil {
//log.NewError(operationID, utils.GetSelfFuncName(), err.Error(), groupID, userID, userMinSeq, minSeq)
@@ -497,16 +510,16 @@ func (db *MsgDatabase) DeleteUserMsgsAndSetMinSeq(ctx context.Context, userID st
if minSeq == 0 {
return nil
}
return db . msgC ache. SetUserMinSeq ( ctx , userID , uint64 ( minSeq ) )
return db . c ache. SetUserMinSeq ( ctx , userID , minSeq )
}
// this is struct for recursion
type delMsgRecursionStruct struct {
minSeq u int32
minSeq int64
delDocIDList [ ] string
}
func ( d * delMsgRecursionStruct ) getSetMinSeq ( ) u int32 {
func ( d * delMsgRecursionStruct ) getSetMinSeq ( ) int64 {
return d . minSeq
}
@@ -514,9 +527,9 @@ func (d *delMsgRecursionStruct) getSetMinSeq() uint32 {
// seq 70
// set minSeq 21
// recursion 删除list并且返回设置的最小seq
func ( db * MsgDatabase ) deleteMsgRecursion ( ctx context . Context , sourceID string , index int64 , delStruct * delMsgRecursionStruct , remainTime int64 ) ( u int32 , error ) {
func ( db * MsgDatabase ) deleteMsgRecursion ( ctx context . Context , sourceID string , index int64 , delStruct * delMsgRecursionStruct , remainTime int64 ) ( int64 , error ) {
// find from oldest list
msgs , err := db . msgModel . GetMsgsByIndex ( ctx , sourceID , index )
msgs , err := db . mgo . GetMsgsByIndex ( ctx , sourceID , index )
if err != nil || msgs . DocID == "" {
if err != nil {
if err == unrelation . ErrMsgListNotExist {
@@ -526,14 +539,14 @@ func (db *MsgDatabase) deleteMsgRecursion(ctx context.Context, sourceID string,
}
}
// 获取报错, 或者获取不到了, 物理删除并且返回seq delMongoMsgsPhysical(delStruct.delDocIDList)
err = db . msgModel . Delete ( ctx , delStruct . delDocIDList )
err = db . mgo . Delete ( ctx , delStruct . delDocIDList )
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))
if len ( msgs . Msg ) > db . msg . GetSingleGocMsgNum ( ) {
if int64 ( len( msgs . Msg ) ) > db . msg . GetSingleGocMsgNum ( ) {
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 ( ) {
@@ -561,11 +574,11 @@ func (db *MsgDatabase) deleteMsgRecursion(ctx context.Context, sourceID string,
msg . SendTime = 0
hasMarkDelFlag = true
} else {
if err := db . msgModel . Delete ( ctx , delStruct . delDocIDList ) ; err != nil {
if err := db . mgo . Delete ( ctx , delStruct . delDocIDList ) ; err != nil {
return 0 , err
}
if hasMarkDelFlag {
if err := db . msgModel . UpdateOneDoc ( ctx , msgs ) ; err != nil {
if err := db . mgo . UpdateOneDoc ( ctx , msgs ) ; err != nil {
return delStruct . getSetMinSeq ( ) , utils . Wrap ( err , "" )
}
}
@@ -578,3 +591,62 @@ func (db *MsgDatabase) deleteMsgRecursion(ctx context.Context, sourceID string,
seq , err := db . deleteMsgRecursion ( ctx , sourceID , index + 1 , delStruct , remainTime )
return seq , utils . Wrap ( err , "deleteMsg failed" )
}
func ( db * MsgDatabase ) GetUserMinMaxSeqInMongoAndCache ( ctx context . Context , userID string ) ( minSeqMongo , maxSeqMongo , minSeqCache , maxSeqCache int64 , err error ) {
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
}
func ( db * MsgDatabase ) GetSuperGroupMinMaxSeqInMongoAndCache ( ctx context . Context , groupID string ) ( minSeqMongo , maxSeqMongo , maxSeqCache int64 , err error ) {
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
}
func ( db * MsgDatabase ) GetMinMaxSeqMongo ( ctx context . Context , sourceID string ) ( minSeqMongo , maxSeqMongo int64 , err error ) {
oldestMsgMongo , err := db . mgo . GetOldestMsg ( ctx , sourceID )
if err != nil {
return 0 , 0 , err
}
msgPb , err := db . unmarshalMsg ( oldestMsgMongo )
if err != nil {
return 0 , 0 , err
}
minSeqMongo = msgPb . Seq
newestMsgMongo , err := db . mgo . GetNewestMsg ( ctx , sourceID )
if err != nil {
return 0 , 0 , err
}
msgPb , err = db . unmarshalMsg ( newestMsgMongo )
if err != nil {
return 0 , 0 , err
}
maxSeqMongo = msgPb . Seq
return
}
func ( db * MsgDatabase ) SetGroupUserMinSeq ( ctx context . Context , groupID , userID string , minSeq int64 ) ( err error ) {
return db . cache . SetGroupUserMinSeq ( ctx , groupID , userID , minSeq )
}
func ( db * MsgDatabase ) SetUserMinSeq ( ctx context . Context , userID string , minSeq int64 ) ( err error ) {
return db . cache . SetUserMinSeq ( ctx , userID , minSeq )
}