package messagemod import ( "91porn-server/common/db" "91porn-server/common/log" "91porn-server/models" "encoding/json" "fmt" "time" "go.mongodb.org/mongo-driver/bson" "go.mongodb.org/mongo-driver/bson/primitive" "go.mongodb.org/mongo-driver/mongo" "go.mongodb.org/mongo-driver/mongo/options" ) var mdb *db.MongoDB const table = models.ChatMessage const NoRedDynamicNumRedisKey = "user-no-red-dynamic-num" // Coll 获取表名 func coll(t *db.MongoTool) *db.MongoTool { if t == nil { return mdb.Coll(table) } return t.Coll(table) } // InitIndex 初始化索引 func initIndex() { many := []mongo.IndexModel{ { Keys: bson.D{{"sendUid", 1}}, }, { Keys: bson.D{{"takeUid", 1}}, }, { Keys: bson.D{{"msgType", 1}}, }, { Keys: bson.D{{"sessionId", 1}}, }, { Keys: bson.D{{"createdAt", -1}}, }, { Keys: bson.D{{"msgType", 1}, {"sessionId", 1}, {"takeUid", 1}, {"isRead", 1}}, }, { Keys: bson.D{{"sessionId", 1}, {"takeUid", 1}, {"isRead", 1}}, }, } _, err := coll(nil).CreateIndex(many) if err != nil { panic(fmt.Sprintf("%s model set index err ==>[%+v]", table, err)) } return } // InsertOne 新增 func InsertOne(mt *db.MongoTool, t *Message) (data primitive.ObjectID, err error) { result, err := coll(mt).InsertOne(&t) if err != nil { log.Error(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "InsertOne", table, "InsertOne", err)) return } byteID, err := json.Marshal(result.InsertedID) if err != nil { log.Error(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "InsertOne", table, "Marshal", err)) return } err = data.UnmarshalJSON(byteID) if err != nil { log.Error(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "InsertOne", table, "UnmarshalJSON", err)) return } return } // 查询总条数 func Count(filter primitive.M) (int64, error) { if count, err := coll(nil).Count(filter); err != nil { log.Error(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "Count", models.Activity, "Count", err), log.Any("filter", filter), ) return 0, err } else { return count, nil } } // DeleteMessages 删除 func DeleteMessages(ids []primitive.ObjectID) (int64, error) { if ids == nil { ids = []primitive.ObjectID{} } cond := bson.M{"_id": bson.M{"$in": ids}} result, err := coll(nil).DeleteMany(cond) if err != nil { log.Error(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "DeleteMessages", table, "DeleteMany", err), log.Any("ids", ids)) return 0, err } return result.DeletedCount, nil } // DeleteMsgBySessionId 根据会话消息删除会话聊天消息 func DeleteMsgBySessionId(sessionId string) (int64, error) { cond := bson.M{"sessionId": sessionId, "msgType": PrivateLetterMsg} result, err := coll(nil).DeleteMany(cond) if err != nil { log.Error(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "DeleteMsgBySessionId", table, "DeleteMany", err), log.Any("sessionId", sessionId)) return 0, err } return result.DeletedCount, nil } // FindOneByFilter 根据条件查询单个信息 func FindOneByFilter(cond bson.M) (data Message, err error) { opts := options.FindOne().SetSort(bson.D{{"createdAt", -1}}) err = coll(nil).FindOne(&data, cond, opts) if err != nil { log.Error(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "FindOneByFilter", table, "FindOne", err), log.Any("cond", cond), ) return } return } // Query 查询列表,app使用 func Query(cond bson.M, pageSize, pageNumber int) (data []MessageApp, hasNext bool, err error) { data = make([]MessageApp, 0) skip := pageSize * (pageNumber - 1) sort := bson.D{{"createdAt", -1}} opts := options.Find().SetSkip(int64(skip)).SetLimit(int64(pageSize) + 1).SetSort(sort) err = coll(nil).Find(&data, cond, opts) if err != nil { log.Error("Query", log.Any("cond", cond), log.Any("skip", skip), log.Any("limit", pageSize), log.E(err)) return } if len(data) > pageSize { hasNext = true data = data[:pageSize] } return } // 查询动态列表,app使用 func QueryDynamics(cond bson.M, pageSize, pageNumber int) (data []MsgDynamicsApp, hasNext bool, err error) { data = make([]MsgDynamicsApp, 0) skip := pageSize * (pageNumber - 1) sort := bson.D{{"createdAt", -1}} opts := options.Find().SetSkip(int64(skip)).SetLimit(int64(pageSize) + 1).SetSort(sort) err = coll(nil).Find(&data, cond, opts) if err != nil { log.Error("QueryDynamics", log.Any("cond", cond), log.Any("skip", skip), log.Any("limit", pageSize), log.E(err)) return } if len(data) > pageSize { hasNext = true data = data[:pageSize] } return } // FindByFilter 条件查询次数,app使用 func FindByFilter(cond bson.M) (data []*MsgDynamicsApp, err error) { err = coll(nil).Find(&data, cond) if err != nil { log.Error("FindByFilter", log.Any("cond", cond), log.E(err)) return } return } // UpdIsRead 更新消息已读状态 func UpdIsRead(t *db.MongoTool, uid uint64, ids []primitive.ObjectID) (err error) { filter := bson.M{"_id": bson.M{"$in": ids}, "isRead": false, "takeUid": uid} set := bson.M{"$set": bson.M{"isRead": true, "updatedAt": time.Now()}} _, err = coll(t).UpdateMany(filter, set) if err != nil { log.Warn(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "UpdIsRead", table, "UpdateMany", err), log.Any("ids", ids), ) return err } return } // UpdIsReadBySessionIdAndUid 根据sessionId更新消息已读状态 func UpdIsReadBySessionIdAndUid(t *db.MongoTool, uid uint64, sessionId string) (err error) { msgTypes := [...]string{string(PrivateLetterMsg), string(OfficialPrivateLetterMsg)} filter := bson.M{"msgType": bson.M{"$in": msgTypes}, "sessionId": sessionId, "takeUid": uid, "isRead": false} set := bson.M{"$set": bson.M{"isRead": true}} _, err = coll(t).UpdateMany(filter, set) if err != nil { log.Warn(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "UpdIsReadBySessionIdAndUid", table, "UpdateMany", err), log.Any("uid", uid), log.Any("sessionId", sessionId)) return err } return } // UpdDynamicsIsReadByUid 根据uid更新动态消息已读状态 func UpdDynamicsIsReadByUid(t *db.MongoTool, uid uint64) (err error) { filter := bson.M{"msgType": bson.M{"$in": GetDynamics()}, "takeUid": uid, "isRead": false} set := bson.M{"$set": bson.M{"isRead": true}} _, err = coll(t).UpdateMany(filter, set) if err != nil { log.Warn(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "UpdDynamicsIsReadByUid", table, "UpdateMany", err), log.Any("uid", uid)) return err } return } // QueryAllDocument 分页查询文档 func QueryAllDocument(filter primitive.M, opts ...*options.FindOptions) ([]*Message, error) { var out []*Message if err := coll(nil).Find(&out, filter, opts...); err != nil { log.Error(fmt.Sprintf("[METHOD-QueryAllDocument]==> Model %s Find fail error:%+v:", table, err), log.Any("filter", filter), ) return nil, err } return out, nil } // CountDocument 查询文档条目数 func CountDocument(filter primitive.M) (int64, error) { count, err := coll(nil).Count(filter) if err != nil { log.Error(fmt.Sprintf("[METHOD-CountDocument]==> Model %s Count fail error:%+v:", table, err), log.Any("filter", filter), ) return 0, err } return count, nil } // RemoveDocument 删除用户私聊信息 func RemoveDocument(filter primitive.M) error { _, err := coll(nil).DeleteOne(filter) if err != nil { log.Error(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "RemoveDocument", table, "DeleteOne", err), log.Any("filter", filter), ) } return err } // UpdOneByFilter 更新 func UpdOneByFilter(t *db.MongoTool, cond bson.M, set bson.M) (err error) { _, err = coll(t).UpdateOne(cond, set) if err != nil { log.Warn(fmt.Sprintf("[METHOD-%s]==> Model %s %s fail error:%+v:", "UpdOneByFilter", table, "UpdateOne", err), log.Any("cond", cond), log.Any("set", cond)) return err } return }