Files
rootandClaude Opus 5 8679200f41 Initial commit
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-15 13:57:10 +08:00

786 lines
21 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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)
}