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

459 lines
13 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 imser
import (
"fmt"
"net/url"
"strconv"
"strings"
"sync"
"time"
"91porn-server/app/appg"
"91porn-server/common"
"91porn-server/common/constant/redisconst"
"91porn-server/common/imclient"
"91porn-server/common/log"
"91porn-server/common/redis"
"91porn-server/models/v/imusermod"
"91porn-server/models/v/usermod"
"github.com/go-redsync/redsync/v4"
)
const (
defaultIMV2DynamicConfigDomain = "https://sdk.api.cdflow.cn"
defaultIMV2SocketURL = "wss://wss.cdflow.cn:32214"
// ensureSDKUser 抢锁失败后,轮询 imUser 表的间隔与最大次数(10 × 200ms = 最多等 2s
imSDKEnsureWaitInterval = 200 * time.Millisecond
imSDKEnsureWaitTries = 10
)
var appTokenCache = struct {
sync.RWMutex
token string
expiresAt time.Time
}{}
type SDKAuthInfo struct {
Enabled bool `json:"enabled"`
ImToken string `json:"imToken"`
Token string `json:"token,omitempty"`
UserID uint64 `json:"userId"`
ImUserID int64 `json:"imUserId"`
Registered bool `json:"registered"`
NewlyRegistered bool `json:"newlyRegistered,omitempty"`
DynamicConfigDomain string `json:"dynamicConfigDomain"`
SocketURL string `json:"socketURL"`
SysType string `json:"sysType,omitempty"`
OS string `json:"os,omitempty"`
OSType string `json:"osType,omitempty"`
}
type SDKPublicConfig struct {
Enable bool `json:"enable"`
BaseURL string `json:"baseUrl"`
SocketURL string `json:"socketUrl"`
MerchantCode string `json:"merchantCode"`
TenantCode string `json:"tenantCode"`
}
func SDKEnabled() bool {
return newSDKClient().Enabled()
}
func GetSDKPublicConfig() SDKPublicConfig {
cfg := sdkConfig()
return SDKPublicConfig{
Enable: newSDKClient().Enabled(),
BaseURL: cfg.BaseURL,
SocketURL: appg.Conf.ImV2.SocketURL,
MerchantCode: cfg.MerchantCode,
TenantCode: cfg.TenantCode,
}
}
func GetSDKAuth(uid uint64) (*SDKAuthInfo, error) {
c := newSDKClient()
if !c.Enabled() {
reason := c.DisabledReason()
log.Error("GetSDKAuth im sdk disabled", log.Any("uid", uid), log.Any("reason", reason))
return nil, fmt.Errorf("im sdk is disabled: %s", reason)
}
// 确保 IM 用户存在(内部以 imUser 表为准 + 并发锁),并得知本次是否新注册
imUserID, newlyRegistered, err := ensureSDKUser(uid)
if err != nil {
return nil, err
}
if err := withAppToken(c, func(token string) error {
return c.SetOnlineStatus(imclient.SetOnlineStatusRequest{
UserID: imUserID,
ShowOnlineStatus: true,
}, token)
}); err != nil {
log.Warn("set im online status visible failed", log.Any("uid", uid), log.Any("imUserId", imUserID), log.E(err))
}
token, err := cachedUserToken(c, imUserID)
if err != nil {
log.Error("get user im token fail", log.Any("uid", uid), log.Any("imUserId", imUserID), log.E(err))
return nil, err
}
return &SDKAuthInfo{
Enabled: true,
ImToken: token,
Token: token,
UserID: uid,
ImUserID: imUserID,
Registered: imUserID > 0,
NewlyRegistered: newlyRegistered,
DynamicConfigDomain: dynamicConfigDomain(),
SocketURL: socketURL(),
OS: "web",
OSType: "web",
}, nil
}
func SendPrivateMessage(sendUID, takeUID uint64, content string, imgURLs []string) error {
c := newSDKClient()
if !c.Enabled() {
return nil
}
senderID, err := EnsureSDKUser(sendUID)
if err != nil {
return err
}
receiverID, err := EnsureSDKUser(takeUID)
if err != nil {
return err
}
if content != "" {
if err = withAppToken(c, func(token string) error {
return c.SendMessage(imclient.SendMessageRequest{
SenderID: senderID,
ReceiverID: receiverID,
Content: content,
MessageType: imclient.MessageTypeText,
}, token)
}); err != nil {
return err
}
}
for _, imgURL := range imgURLs {
if imgURL == "" {
continue
}
if err = withAppToken(c, func(token string) error {
return c.SendMessage(imclient.SendMessageRequest{
SenderID: senderID,
ReceiverID: receiverID,
Content: "",
MessageType: imclient.MessageTypeImage,
Attachment: &imclient.Attachment{
URL: imgURL,
},
}, token)
}); err != nil {
return err
}
}
return nil
}
func SyncSDKBaseInfo(uid uint64, name, portrait, signature *string) error {
c := newSDKClient()
if !c.Enabled() {
return nil
}
imUserID, err := EnsureSDKUser(uid)
if err != nil {
return err
}
req := imclient.UpdateBaseInfoRequest{
UserID: imUserID,
}
if name != nil {
req.Name = *name
}
if portrait != nil {
req.ImgURL = common.BindUrl(appg.Conf.URL.OriginUrl, *portrait)
}
if signature != nil {
req.Signature = *signature
}
if req.Name == "" && req.ImgURL == "" && req.Signature == "" {
return nil
}
return withAppToken(c, func(token string) error {
return c.UpdateBaseInfo(req, token)
})
}
// EnsureSDKUser 确保 IM 用户存在并返回 imUserID。
func EnsureSDKUser(uid uint64) (int64, error) {
imUserID, _, err := ensureSDKUser(uid)
return imUserID, err
}
// ensureSDKUser 确保 IM 用户存在并返回 imUserIDnewlyRegistered 表示本次是否真正发起了 SDK 注册
// (命中 imUser 表或回填历史用户均为 false)。以 imUser 表为准判断是否已注册。
func ensureSDKUser(uid uint64) (int64, bool, error) {
if uid == 0 {
log.Warn("EnsureSDKUser uid is zero")
return 0, false, fmt.Errorf("uid is empty")
}
// fast path:以 imUser 表为准,已注册直接返回
if mapping, _ := imusermod.FindByUIDQuiet(uid); mapping.IMUserID > 0 {
return mapping.IMUserID, false, nil
}
// 未注册 → 抢锁串行化,防止并发重复注册。只抢一次、不自旋。
mu := redis.BuildLock(appg.Redis, redisconst.IMSDKEnsureLockKey(uid),
redsync.WithExpiry(10*time.Second),
redsync.WithTries(1),
)
if err := mu.Lock(); err != nil {
// 没抢到锁 → 已有别的请求在注册:每 200ms 轮询一次 imUser 表,最多等 2s
for i := 0; i < imSDKEnsureWaitTries; i++ {
time.Sleep(imSDKEnsureWaitInterval)
if mapping, _ := imusermod.FindByUIDQuiet(uid); mapping.IMUserID > 0 {
return mapping.IMUserID, false, nil
}
}
return 0, false, fmt.Errorf("ensure sdk user busy: wait imUser timeout, uid=%d", uid)
}
defer func() { _, _ = mu.Unlock() }()
// 双重检查:等锁期间可能已被其他请求注册完,避免重复注册
if mapping, _ := imusermod.FindByUIDQuiet(uid); mapping.IMUserID > 0 {
return mapping.IMUserID, false, nil
}
u, err := usermod.FindUserByUID(uid)
if err != nil {
log.Warn(fmt.Sprintf("EnsureSDKUser usermod.FindUserByUID fail uid=%d", uid), log.E(err))
return 0, false, err
}
if u == nil {
log.Warn(fmt.Sprintf("EnsureSDKUser user is nil uid=%d", uid))
return 0, false, fmt.Errorf("user %d not found", uid)
}
// 一切以 imUser 表为准:表里无映射即视为未注册,直接走 SDK 注册(不信任 user.ImUserID,不回填)
c := newSDKClient()
thirdPartyID := strconv.FormatUint(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(uid),
Nickname: u.Name,
Avatar: u.Portrait,
}, token)
return registerErr
})
if err != nil {
log.Warn(fmt.Sprintf("EnsureSDKUser c.Register fail uid=%d", uid), log.E(err))
return 0, false, err
}
if imUserID <= 0 {
log.Warn(fmt.Sprintf("EnsureSDKUser im register returned empty user id uid=%d", uid), log.E(err))
return 0, false, fmt.Errorf("im register returned empty user id")
}
if err = imusermod.UpsertByUID(uid, imUserID, thirdPartyID); err != nil {
log.Warn(fmt.Sprintf("EnsureSDKUser upsert fail uid=%d", uid), log.E(err))
return 0, false, err
}
return imUserID, true, nil
}
// InitAppToken 启动时异步取一次 appToken(不阻断启动),并每小时刷新一次,写入内存+redis。
func InitAppToken() {
if !SDKEnabled() {
return
}
common.Go(func() {
if err := refreshAppToken(); err != nil {
log.Warn("init app token failed", log.E(err))
}
})
common.Go(func() {
ticker := time.NewTicker(time.Hour)
defer ticker.Stop()
for range ticker.C {
if err := refreshAppToken(); err != nil {
log.Warn("refresh app token failed", log.E(err))
}
}
})
}
// refreshAppToken 从 SDK 取最新 appToken,写入内存 + redis。
func refreshAppToken() error {
c := newSDKClient()
if !c.Enabled() {
return nil
}
token, err := c.AppToken()
if err != nil {
return err
}
if token == "" {
return fmt.Errorf("im app token empty")
}
ttl := sdkTokenTTL()
cacheSeconds := ttl - 60
if cacheSeconds < 60 {
cacheSeconds = ttl
}
setToken(token, time.Duration(cacheSeconds)*time.Second)
return nil
}
// setToken 写入内存 + redis。
func setToken(token string, ttl time.Duration) {
appTokenCache.Lock()
appTokenCache.token = token
appTokenCache.expiresAt = time.Now().Add(ttl)
appTokenCache.Unlock()
if err := appg.Redis.Set(redisconst.IMSDKAppTokenKey, token, ttl); err != nil {
log.Warn("set app token to redis fail", log.E(err))
}
}
// getToken 取 appToken:优先内存(读锁),取不到再取 redis(命中则回种内存)。
func getToken() string {
appTokenCache.RLock()
token := appTokenCache.token
appTokenCache.RUnlock()
if token != "" {
return token
}
if v, err := appg.Redis.Get(redisconst.IMSDKAppTokenKey); err == nil && v != nil && *v != "" {
appTokenCache.Lock()
appTokenCache.token = *v
appTokenCache.Unlock()
return *v
}
return ""
}
func appToken(c *imclient.Client) (string, error) {
if token := getToken(); token != "" {
return token, nil
}
// 内存和 redis 都没有(启动初期或刷新失败)→ 兜底同步取一次
if err := refreshAppToken(); err != nil {
return "", err
}
if token := getToken(); token != "" {
return token, nil
}
return "", fmt.Errorf("im app token unavailable")
}
func invalidateAppToken() {
appTokenCache.Lock()
appTokenCache.token = ""
appTokenCache.expiresAt = time.Time{}
appTokenCache.Unlock()
_, _ = appg.Redis.Del(redisconst.IMSDKAppTokenKey)
}
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")
}
// cachedUserToken 获取 IM 用户 token,缓存 12 小时,避免每次鉴权都请求 SDK。
// 缓存按 imUserID 维度,Redis 故障时降级为直接请求,不阻断鉴权。
func cachedUserToken(c *imclient.Client, imUserID int64) (string, error) {
key := redisconst.IMSDKUserTokenKey(imUserID)
if cached, err := appg.Redis.Get(key); err != nil {
log.Warn("get im user token cache fail", log.Any("imUserId", imUserID), log.E(err))
} else if cached != nil && *cached != "" {
return *cached, nil
}
token, err := c.UserToken(imUserID)
if err != nil {
return "", err
}
if token != "" {
if err = appg.Redis.Set(key, token, redisconst.IMSDKUserTokenExpire); err != nil {
log.Warn("set im user token cache fail", log.Any("imUserId", imUserID), log.E(err))
}
}
return token, nil
}
func sdkTokenTTL() int64 {
ttl := sdkConfig().TokenTTL
if ttl <= 0 {
ttl = 86400
}
return ttl
}
func newSDKClient() *imclient.Client {
return imclient.New(sdkConfig())
}
func sdkConfig() imclient.Config {
cfg := appg.Conf.ImV2
return 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: 86400,
}
}
func dynamicConfigDomain() string {
if domain := strings.TrimRight(appg.Conf.ImV2.DynamicConfigDomain, "/"); domain != "" {
return domain
}
rawBaseURL := strings.TrimSpace(appg.Conf.ImV2.BaseURL)
if rawBaseURL != "" {
if u, err := url.Parse(rawBaseURL); err == nil && u.Scheme != "" && u.Host != "" {
return u.Scheme + "://" + u.Host
}
}
return defaultIMV2DynamicConfigDomain
}
func socketURL() string {
if socket := strings.TrimSpace(appg.Conf.ImV2.SocketURL); socket != "" {
return socket
}
return defaultIMV2SocketURL
}
func sdkPassword(uid uint64) string {
return fmt.Sprintf("hjll%014d", uid%100000000000000)
}
func LogSDKSendError(sendUID, takeUID uint64, err error) {
if err == nil {
return
}
log.Warn("im sdk send private message error", log.Any("sendUID", sendUID), log.Any("takeUID", takeUID), log.E(err))
}