package srv_im import ( "encoding/json" "fmt" "strconv" "strings" "sync" "time" "91porn-server/common/imclient" "91porn-server/common/stderr" "91porn-server/models/v/imusermod" "91porn-server/models/v/sessionmod" "91porn-server/models/v/usermod" "91porn-server/skd/skdg" "91porn-server/web/webg" "go.mongodb.org/mongo-driver/bson" "go.mongodb.org/mongo-driver/bson/primitive" "go.mongodb.org/mongo-driver/mongo/options" ) const ( maxIMOnlineStatusBatchSize = 100 defaultIMTokenTTL = int64(86400) imUserSyncPageSize = int64(500) ) var webAppTokenCache = struct { sync.Mutex token string expiresAt time.Time }{} type UserListReq struct { UID *uint64 `json:"uid"` HasIM *bool `json:"hasIm"` IsUp *bool `json:"isUp"` PageNum int64 `json:"pageNum"` PageSize int64 `json:"pageSize"` } type UserInfo struct { UID uint64 `json:"uid"` ImUserID int64 `json:"imUserId"` NickName string `json:"nickName"` Avatar string `json:"avatar"` IsUp bool `json:"isUp"` } type UserListResp struct { Total int64 `json:"total"` List []UserInfo `json:"list"` } type SyncUsersReq struct { UIDs []uint64 `json:"uids" binding:"required"` } type SyncUsersResp struct { Total int `json:"total"` Success int `json:"success"` Skipped int `json:"skipped"` Failed int `json:"failed"` } type FriendAddDirectReq struct { UID uint64 `json:"uid" binding:"required"` PeerUID uint64 `json:"peerUid" binding:"required"` } type OnlineStatusReq struct { UIDs []uint64 `json:"uids" binding:"required"` } type OnlineStatusInfo struct { UID uint64 `json:"uid"` ImUserID int64 `json:"imUserId"` Online bool `json:"online"` } type PassthroughReq struct { UIDs []uint64 `json:"uids"` UserIDs []uint64 `json:"userIds"` PassthroughType string `json:"passthroughType"` EventType string `json:"eventType"` Content string `json:"content"` ExtInfo string `json:"extInfo"` Payload map[string]interface{} `json:"payload"` ChannelType string `json:"channelType"` SenderUID uint64 `json:"senderUid"` AllApp bool `json:"allApp"` OnlineOnly bool `json:"onlineOnly"` } type PassthroughResp struct { TargetCount int `json:"targetCount"` Result *imclient.PassthroughResult `json:"result,omitempty"` } type SendMessageReq struct { UID uint64 `json:"uid" binding:"required"` PeerUID uint64 `json:"peerUid" binding:"required"` Content string `json:"content" binding:"required"` } type SendMessageResp struct { UID uint64 `json:"uid"` ImUserID int64 `json:"imUserId"` PeerUID uint64 `json:"peerUid"` PeerImUserID int64 `json:"peerImUserId"` Content string `json:"content"` MessageType int `json:"messageType"` } type SyncDialogUsersStats struct { Enabled bool `json:"enabled"` Total int `json:"total"` Skipped int `json:"skipped"` Success int `json:"success"` Failed int `json:"failed"` FriendSuccess int `json:"friendSuccess"` FriendFailed int `json:"friendFailed"` } type imFriendPair struct { UID uint64 PeerUID uint64 } func ListUsers(req UserListReq) (UserListResp, stderr.Code, string) { if req.PageNum <= 0 { req.PageNum = 1 } if req.PageSize <= 0 || req.PageSize > 100 { req.PageSize = 20 } filter := buildUserListFilter(req) skip := (req.PageNum - 1) * req.PageSize opts := options.Find(). SetSort(bson.M{"uid": -1}). SetSkip(skip). SetLimit(req.PageSize) var total int64 users, _, err := usermod.FetchList(filter, opts, &total) if err != nil { return UserListResp{}, stderr.Failure, err.Error() } resp := UserListResp{Total: total, List: make([]UserInfo, 0, len(users))} for _, user := range users { resp.List = append(resp.List, userInfoFromUser(user)) } return resp, stderr.Success, "" } func SyncUsers(req SyncUsersReq) (SyncUsersResp, stderr.Code, string) { uids := normalizeUIDs(req.UIDs) resp := SyncUsersResp{Total: len(uids)} if len(uids) == 0 { return resp, stderr.ErrParamError, "uids不能为空" } c := newWebSDKClient() if !c.Enabled() { return resp, stderr.Failure, "IM 未配置" } for _, uid := range uids { user, err := usermod.FindUserByUID(uid) if err != nil || user == nil || user.UID == 0 { resp.Failed++ continue } if imusermod.IMUserIDByUID(user.UID) > 0 { resp.Skipped++ continue } if _, err = ensureUserRegistered(c, user); err != nil { resp.Failed++ continue } resp.Success++ } return resp, stderr.Success, "" } func AddFriendDirect(req FriendAddDirectReq) (stderr.Code, string) { if req.UID == 0 || req.PeerUID == 0 || req.UID == req.PeerUID { return stderr.ErrParamError, "uid/peerUid错误" } c := newWebSDKClient() userIMID, peerIMID, err := ensurePairRegistered(c, req.UID, req.PeerUID) if err != nil { return stderr.Failure, err.Error() } if err = ensureFriendDirect(c, userIMID, peerIMID); err != nil { return stderr.Failure, err.Error() } if err = ensureFriendDirect(c, peerIMID, userIMID); err != nil { return stderr.Failure, err.Error() } return stderr.Success, "" } func BatchOnlineStatus(req OnlineStatusReq) ([]OnlineStatusInfo, stderr.Code, string) { uids := normalizeUIDs(req.UIDs) if len(uids) == 0 { return []OnlineStatusInfo{}, stderr.Success, "" } mappings, err := imusermod.FindByUIDsQuiet(uids) if err != nil { return nil, stderr.Failure, err.Error() } imIDs := make([]int64, 0, len(mappings)) byIMID := make(map[int64]uint64, len(mappings)) seen := make(map[int64]struct{}, len(mappings)) for _, m := range mappings { if m.IMUserID <= 0 { continue } if _, ok := seen[m.IMUserID]; ok { continue } seen[m.IMUserID] = struct{}{} imIDs = append(imIDs, m.IMUserID) byIMID[m.IMUserID] = m.UID } if len(imIDs) == 0 { return []OnlineStatusInfo{}, stderr.Success, "" } c := newWebSDKClient() statuses := make([]imclient.OnlineStatus, 0, len(imIDs)) for _, batch := range chunkInt64s(imIDs, maxIMOnlineStatusBatchSize) { var items []imclient.OnlineStatus err := withAppToken(c, func(token string) error { var batchErr error items, batchErr = c.BatchOnlineStatus(imclient.BatchOnlineStatusRequest{UserIDs: batch}, token) return batchErr }) if err != nil { return nil, stderr.Failure, err.Error() } statuses = append(statuses, items...) } resp := make([]OnlineStatusInfo, 0, len(statuses)) for _, item := range statuses { resp = append(resp, OnlineStatusInfo{ UID: byIMID[item.UserID], ImUserID: item.UserID, Online: item.Online, }) } return resp, stderr.Success, "" } func SendPassthrough(req PassthroughReq) (PassthroughResp, stderr.Code, string) { passthroughType := strings.TrimSpace(req.PassthroughType) if passthroughType == "" { passthroughType = strings.TrimSpace(req.EventType) } if passthroughType == "" { return PassthroughResp{}, stderr.ErrParamError, "passthroughType不能为空" } c := newWebSDKClient() content := strings.TrimSpace(req.Content) extInfo := strings.TrimSpace(req.ExtInfo) if req.Payload != nil { payloadBytes, err := json.Marshal(req.Payload) if err != nil { return PassthroughResp{}, stderr.ErrParamError, err.Error() } if content == "" { content = string(payloadBytes) } if extInfo == "" { extInfo = string(payloadBytes) } } if extInfo == "" { extInfo = "{}" } senderID, err := resolveSenderIMID(c, req.SenderUID) if err != nil { return PassthroughResp{}, stderr.Failure, err.Error() } if req.AllApp { var result *imclient.PassthroughResult err := withAppToken(c, func(token string) error { var sendErr error result, sendErr = c.SendAppPassthrough(imclient.AppPassthroughRequest{ SenderID: senderID, PassthroughType: passthroughType, Content: content, ExtInfo: extInfo, ChannelType: req.ChannelType, }, token) return sendErr }) if err != nil { return PassthroughResp{}, stderr.Failure, err.Error() } return PassthroughResp{Result: result}, stderr.Success, "" } uids := normalizeUIDs(append(req.UIDs, req.UserIDs...)) if len(uids) == 0 { return PassthroughResp{}, stderr.ErrParamError, "uids不能为空;全应用透传需显式传 allApp=true" } imIDs, err := loadRegisteredIMIDs(uids) if err != nil { return PassthroughResp{}, stderr.Failure, err.Error() } if len(imIDs) == 0 { return PassthroughResp{}, stderr.ErrParamError, "没有匹配到可发送透传的用户" } var result *imclient.PassthroughResult err = withAppToken(c, func(token string) error { var sendErr error result, sendErr = c.SendOnlinePassthrough(imclient.OnlinePassthroughRequest{ SenderID: senderID, ReceiverIDSet: imIDs, PassthroughType: passthroughType, Content: content, ExtInfo: extInfo, ChannelType: req.ChannelType, }, token) return sendErr }) if err != nil { return PassthroughResp{}, stderr.Failure, err.Error() } return PassthroughResp{TargetCount: len(imIDs), Result: result}, stderr.Success, "" } func SendMessage(req SendMessageReq) (SendMessageResp, stderr.Code, string) { content := strings.TrimSpace(req.Content) if content == "" { return SendMessageResp{}, stderr.ErrParamError, "content不能为空" } if req.UID == 0 || req.PeerUID == 0 || req.UID == req.PeerUID { return SendMessageResp{}, stderr.ErrParamError, "uid/peerUid错误" } c := newWebSDKClient() userIMID, peerIMID, err := ensurePairRegistered(c, req.UID, req.PeerUID) if err != nil { return SendMessageResp{}, stderr.Failure, err.Error() } if err = ensureFriendDirect(c, userIMID, peerIMID); err != nil { return SendMessageResp{}, stderr.Failure, err.Error() } if err = ensureFriendDirect(c, peerIMID, userIMID); err != nil { return SendMessageResp{}, stderr.Failure, err.Error() } messageType := imclient.MessageTypeText if err = withAppToken(c, func(token string) error { return c.SendMessage(imclient.SendMessageRequest{ SenderID: userIMID, ReceiverID: peerIMID, Content: content, MessageType: messageType, }, token) }); err != nil { return SendMessageResp{}, stderr.Failure, err.Error() } return SendMessageResp{ UID: req.UID, ImUserID: userIMID, PeerUID: req.PeerUID, PeerImUserID: peerIMID, Content: content, MessageType: messageType, }, stderr.Success, "" } func SyncDialogUsers() (SyncDialogUsersStats, error) { c := newWebSDKClient() if !c.Enabled() { return SyncDialogUsersStats{Enabled: false}, nil } syncIDs, pairs, err := CollectDialogUserIDs() if err != nil { return SyncDialogUsersStats{Enabled: true}, err } stats := SyncDialogUsersStats{Enabled: true, Total: len(syncIDs)} imIDs := make(map[uint64]int64, len(syncIDs)) for uid := range syncIDs { user, err := usermod.FindUserByUID(uid) if err != nil || user == nil || user.UID == 0 { stats.Failed++ continue } if imid := imusermod.IMUserIDByUID(user.UID); imid > 0 { imIDs[user.UID] = imid stats.Skipped++ continue } imUserID, err := ensureUserRegistered(c, user) if err != nil { stats.Failed++ continue } imIDs[user.UID] = imUserID stats.Success++ time.Sleep(100 * time.Millisecond) } for _, pair := range pairs { userIMID := imIDs[pair.UID] peerIMID := imIDs[pair.PeerUID] if userIMID <= 0 || peerIMID <= 0 { stats.FriendFailed++ continue } if err := ensureFriendDirect(c, userIMID, peerIMID); err != nil { stats.FriendFailed++ continue } stats.FriendSuccess++ time.Sleep(100 * time.Millisecond) } return stats, nil } func CollectDialogUserIDs() (map[uint64]struct{}, []imFriendPair, error) { syncIDs := make(map[uint64]struct{}) pairSet := make(map[string]struct{}) pairs := make([]imFriendPair, 0) var lastID primitive.ObjectID for { filter := bson.M{} if !lastID.IsZero() { filter["_id"] = bson.M{"$gt": lastID} } opts := options.Find(). SetProjection(bson.M{"_id": 1, "sendUid": 1, "takeUid": 1}). SetSort(bson.M{"_id": 1}). SetLimit(imUserSyncPageSize) sessions, err := sessionmod.FindManyByFilter(filter, opts) if err != nil { return nil, nil, err } if len(sessions) == 0 { break } for _, session := range sessions { if session.SendUid > 0 { syncIDs[session.SendUid] = struct{}{} } if session.TakeUid > 0 { syncIDs[session.TakeUid] = struct{}{} } addIMFriendPair(pairSet, &pairs, session.SendUid, session.TakeUid) } lastID = sessions[len(sessions)-1].ID if int64(len(sessions)) < imUserSyncPageSize { break } } return syncIDs, pairs, nil } func buildUserListFilter(req UserListReq) bson.M { filter := bson.M{} andFilters := bson.A{} if req.UID != nil && *req.UID > 0 { filter["uid"] = *req.UID } if req.HasIM != nil { if *req.HasIM { filter["imUserId"] = bson.M{"$gt": 0} } else { andFilters = append(andFilters, bson.M{"$or": bson.A{ bson.M{"imUserId": bson.M{"$exists": false}}, bson.M{"imUserId": bson.M{"$lte": 0}}, }}) } } if req.IsUp != nil { if *req.IsUp { andFilters = append(andFilters, creatorFilter()) } else { andFilters = append(andFilters, bson.M{"$nor": bson.A{creatorFilter()}}) } } if len(andFilters) > 0 { filter["$and"] = andFilters } return filter } func creatorFilter() bson.M { return bson.M{"$or": bson.A{ bson.M{"originalUp": true}, bson.M{"officialCert": true}, bson.M{"superUser": true}, bson.M{"merchantUser": bson.M{"$gt": 0}}, bson.M{"vidUploadCount": bson.M{"$gt": 0}}, bson.M{"coverUploadCount": bson.M{"$gt": 0}}, bson.M{"upTag": bson.M{"$ne": ""}}, }} } func userInfoFromUser(user *usermod.User) UserInfo { if user == nil { return UserInfo{} } return UserInfo{ UID: user.UID, ImUserID: imusermod.IMUserIDByUID(user.UID), NickName: user.Name, Avatar: user.Portrait, IsUp: isCreator(user), } } func isCreator(user *usermod.User) bool { return user != nil && (user.OfficialCert || user.SuperUser || user.MerchantUser > 0 || user.VidUploadCount > 0 || user.CoverUploadCount > 0 || user.UpTag != "") } func ensurePairRegistered(c *imclient.Client, uid, peerUID uint64) (int64, int64, error) { user, err := usermod.FindUserByUID(uid) if err != nil { return 0, 0, err } peer, err := usermod.FindUserByUID(peerUID) if err != nil { return 0, 0, err } userIMID, err := ensureUserRegistered(c, user) if err != nil { return 0, 0, fmt.Errorf("sync user im id failed: %w", err) } peerIMID, err := ensureUserRegistered(c, peer) if err != nil { return 0, 0, fmt.Errorf("sync peer im id failed: %w", err) } return userIMID, peerIMID, nil } func ensureUserRegistered(c *imclient.Client, user *usermod.User) (int64, error) { if user == nil || user.UID == 0 { return 0, fmt.Errorf("user is empty") } // 以 imusermod 为准:有映射即已注册 if mapping, err := imusermod.FindByUIDQuiet(user.UID); err == nil && mapping.IMUserID > 0 { return mapping.IMUserID, nil } thirdPartyID := strconv.FormatUint(user.UID, 10) var imUserID int64 err := withAppToken(c, func(token string) error { var registerErr error imUserID, registerErr = c.Register(imclient.RegisterRequest{ ThirdPartyID: thirdPartyID, Password: sdkPassword(user.UID), Nickname: user.Name, Avatar: user.Portrait, }, token) return registerErr }) if err != nil { return 0, err } if imUserID <= 0 { return 0, fmt.Errorf("im register returned empty user id") } if err = imusermod.UpsertByUID(user.UID, imUserID, thirdPartyID); err != nil { return 0, err } return imUserID, nil } func resolveSenderIMID(c *imclient.Client, senderUID uint64) (int64, error) { if senderUID == 0 { return 0, nil } user, err := usermod.FindUserByUID(senderUID) if err != nil { return 0, err } return ensureUserRegistered(c, user) } func loadRegisteredIMIDs(uids []uint64) ([]int64, error) { mappings, err := imusermod.FindByUIDsQuiet(normalizeUIDs(uids)) if err != nil { return nil, err } imIDs := make([]int64, 0, len(mappings)) seen := make(map[int64]struct{}, len(mappings)) for _, m := range mappings { if m.IMUserID <= 0 { continue } if _, ok := seen[m.IMUserID]; ok { continue } seen[m.IMUserID] = struct{}{} imIDs = append(imIDs, m.IMUserID) } return imIDs, nil } func ensureFriendDirect(c *imclient.Client, userIMID, friendIMID int64) error { if userIMID <= 0 || friendIMID <= 0 || userIMID == friendIMID { return nil } err := withAppToken(c, func(token string) error { return c.DirectAddFriend(imclient.DirectAddFriendRequest{ UserID: userIMID, FriendID: friendIMID, Archive: true, }, token) }) if err == nil || isDuplicateFriendError(err) { return nil } return err } func isDuplicateFriendError(err error) bool { if err == nil { return false } msg := strings.ToLower(err.Error()) return strings.Contains(msg, "already") || strings.Contains(msg, "exist") || strings.Contains(msg, "重复") || strings.Contains(msg, "已是好友") || strings.Contains(msg, "好友关系") } func addIMFriendPair(pairSet map[string]struct{}, pairs *[]imFriendPair, uid, peerUID uint64) { addDirectedIMFriendPair(pairSet, pairs, uid, peerUID) addDirectedIMFriendPair(pairSet, pairs, peerUID, uid) } func addDirectedIMFriendPair(pairSet map[string]struct{}, pairs *[]imFriendPair, uid, peerUID uint64) { if uid == 0 || peerUID == 0 || uid == peerUID { return } key := fmt.Sprintf("%d:%d", uid, peerUID) if _, ok := pairSet[key]; ok { return } pairSet[key] = struct{}{} *pairs = append(*pairs, imFriendPair{UID: uid, PeerUID: peerUID}) } func normalizeUIDs(uids []uint64) []uint64 { seen := make(map[uint64]struct{}, len(uids)) normalized := make([]uint64, 0, len(uids)) for _, uid := range uids { if uid == 0 { continue } if _, ok := seen[uid]; ok { continue } seen[uid] = struct{}{} normalized = append(normalized, uid) } return normalized } func chunkInt64s(ids []int64, size int) [][]int64 { if size <= 0 { size = maxIMOnlineStatusBatchSize } if len(ids) == 0 { return nil } chunks := make([][]int64, 0, (len(ids)+size-1)/size) for start := 0; start < len(ids); start += size { end := start + size if end > len(ids) { end = len(ids) } chunks = append(chunks, ids[start:end]) } return chunks } // srvImV2Cfg 拿当前进程可用的 ImV2 配置: // skd cron 调用 SendAdNotify 时 webg.Conf 为 nil,必须用 skdg.Conf; // web 服务调用反之。两侧 imv2 段必须配同样的内容。 type srvImV2 struct { BaseURL string DynamicConfigDomain string SocketURL string MerchantCode string TenantCode string AppKey string ClientID string ClientSecret string } func srvImV2Cfg() srvImV2 { if skdg.Conf != nil { c := skdg.Conf.ImV2 if c.BaseURL != "" || c.AppKey != "" { return srvImV2{ BaseURL: c.BaseURL, DynamicConfigDomain: c.DynamicConfigDomain, SocketURL: c.SocketURL, MerchantCode: c.MerchantCode, TenantCode: c.TenantCode, AppKey: c.AppKey, ClientID: c.ClientID, ClientSecret: c.ClientSecret, } } } if webg.Conf != nil { c := webg.Conf.ImV2 return srvImV2{ BaseURL: c.BaseURL, DynamicConfigDomain: c.DynamicConfigDomain, SocketURL: c.SocketURL, MerchantCode: c.MerchantCode, TenantCode: c.TenantCode, AppKey: c.AppKey, ClientID: c.ClientID, ClientSecret: c.ClientSecret, } } return srvImV2{} } func newWebSDKClient() *imclient.Client { cfg := srvImV2Cfg() return imclient.New(imclient.Config{ Enable: true, BaseURL: cfg.BaseURL, MerchantCode: cfg.MerchantCode, TenantCode: cfg.TenantCode, AppKey: cfg.AppKey, ClientID: cfg.ClientID, ClientSecret: cfg.ClientSecret, SignKey: imclient.DefaultSignKey, AESKey: imclient.DefaultAESKey, EnableSign: true, EncryptTimestamp: true, TokenTTL: defaultIMTokenTTL, }) } func appToken(c *imclient.Client) (string, error) { if !c.Enabled() { return "", fmt.Errorf("IM 未配置") } webAppTokenCache.Lock() defer webAppTokenCache.Unlock() if webAppTokenCache.token != "" && time.Now().Before(webAppTokenCache.expiresAt) { return webAppTokenCache.token, nil } token, err := c.AppToken() if err != nil { return "", err } webAppTokenCache.token = token cacheSeconds := defaultIMTokenTTL - 60 if cacheSeconds < 60 { cacheSeconds = defaultIMTokenTTL } webAppTokenCache.expiresAt = time.Now().Add(time.Duration(cacheSeconds) * time.Second) return token, nil } func invalidateAppToken() { webAppTokenCache.Lock() defer webAppTokenCache.Unlock() webAppTokenCache.token = "" webAppTokenCache.expiresAt = time.Time{} } func withAppToken(c *imclient.Client, fn func(token string) error) error { for attempt := 0; attempt < 2; attempt++ { token, err := appToken(c) if err != nil { return err } err = fn(token) if err == nil { return nil } if attempt == 0 && imclient.IsSessionExpired(err) { invalidateAppToken() continue } return err } return fmt.Errorf("im app token retry exhausted") } func sdkPassword(uid uint64) string { return fmt.Sprintf("hjll%014d", uid%100000000000000) }