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

177 lines
4.9 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 job
import (
"encoding/json"
"fmt"
"time"
"91porn-server/common/log"
"91porn-server/common/redis"
"91porn-server/skd/skdg"
"91porn-server/services/srv_im"
"91porn-server/web/webg"
)
const (
phaseQueued = "queued"
phaseScanningDialogs = "scanning_chat_sessions"
phaseRegisteringUsers = "registering_users"
phaseAddingFriends = "adding_friends"
phaseDone = "done"
syncIMUsersTriggerKey = "im:sync_users:trigger"
syncIMUsersLockKey = "im:sync_users:lock"
syncIMUsersStatusKey = "im:sync_users:status"
syncIMUsersTriggerTTL = 10 * time.Minute
syncIMUsersLockTTL = 1 * time.Hour
)
// SyncIMUsersTaskStatus 旧私聊会话用户同步任务状态。
// 状态完全落在 Redis 里——web 进程写 trigger、读 statusskd 进程消费 trigger、执行任务、写 status。
type SyncIMUsersTaskStatus struct {
Running bool `json:"running"`
Phase string `json:"phase,omitempty"`
StartedAt *time.Time `json:"startedAt,omitempty"`
FinishedAt *time.Time `json:"finishedAt,omitempty"`
LastStats srv_im.SyncDialogUsersStats `json:"lastStats"`
LastError string `json:"lastError,omitempty"`
}
// imSyncRedis 拿当前进程可用的 Redis 客户端:
// skd 主进程用 skdg.Redisweb 主进程用 webg.Redis。
// 两端写的是同一个 Redis 实例,靠 key 协议互通。
func imSyncRedis() *redis.Client {
if skdg.Redis != nil {
return skdg.Redis
}
if webg.Redis != nil {
return webg.Redis
}
return nil
}
// RequestSyncIMUsers 由 web admin 接口调用:写 trigger key,由 skd 接管执行。
// 返回 accepted=true 表示成功排队;false 表示已有任务在跑或在排队。
func RequestSyncIMUsers() (accepted bool, status SyncIMUsersTaskStatus) {
status = GetSyncIMUsersStatus()
if status.Running {
return false, status
}
r := imSyncRedis()
if r == nil {
status.LastError = "redis unavailable"
return false, status
}
if existing, _ := r.Get(syncIMUsersTriggerKey); existing != nil && *existing != "" {
// 已经有 trigger 在排队,等 skd 消费
return false, status
}
if err := r.Set(syncIMUsersTriggerKey, "1", syncIMUsersTriggerTTL); err != nil {
status.LastError = err.Error()
return false, status
}
now := time.Now()
status = SyncIMUsersTaskStatus{
Phase: phaseQueued,
StartedAt: &now,
LastStats: srv_im.SyncDialogUsersStats{Enabled: true},
}
saveSyncIMUsersStatus(status)
return true, status
}
// GetSyncIMUsersStatus 读 Redis 里最新状态,任何进程都能调
func GetSyncIMUsersStatus() SyncIMUsersTaskStatus {
r := imSyncRedis()
if r == nil {
return SyncIMUsersTaskStatus{}
}
raw, err := r.Get(syncIMUsersStatusKey)
if err != nil || raw == nil || *raw == "" {
return SyncIMUsersTaskStatus{}
}
var status SyncIMUsersTaskStatus
if err := json.Unmarshal([]byte(*raw), &status); err != nil {
return SyncIMUsersTaskStatus{}
}
return status
}
// RunPendingSyncIMUsers 由 skd cron 定时调用:检测 trigger 并接管执行同步
func RunPendingSyncIMUsers() {
r := imSyncRedis()
if r == nil {
return
}
// 没有 trigger 就直接退
trig, _ := r.Get(syncIMUsersTriggerKey)
if trig == nil || *trig == "" {
return
}
// 加锁防止 skd 多副本同时跑
ok, err := r.Setnx_NewOK(syncIMUsersLockKey, 1, syncIMUsersLockTTL)
if err != nil || !ok {
return
}
defer func() { _, _ = r.Del(syncIMUsersLockKey) }()
// 消费 trigger,避免下个 tick 重复拉起
_, _ = r.Del(syncIMUsersTriggerKey)
now := time.Now()
status := SyncIMUsersTaskStatus{
Running: true,
Phase: phaseScanningDialogs,
StartedAt: &now,
LastStats: srv_im.SyncDialogUsersStats{Enabled: true},
}
saveSyncIMUsersStatus(status)
stats, runErr := runSyncIMUsersOnce(&status)
finished := time.Now()
status.Running = false
status.Phase = phaseDone
status.FinishedAt = &finished
status.LastStats = stats
if runErr != nil {
status.LastError = runErr.Error()
log.Error("RunPendingSyncIMUsers failed", log.E(runErr), log.Any("stats", stats))
} else {
status.LastError = ""
log.Info("RunPendingSyncIMUsers done", log.Any("stats", stats))
}
saveSyncIMUsersStatus(status)
}
func runSyncIMUsersOnce(status *SyncIMUsersTaskStatus) (stats srv_im.SyncDialogUsersStats, err error) {
defer func() {
if rec := recover(); rec != nil {
err = fmt.Errorf("sync panic: %v", rec)
}
}()
status.Phase = phaseRegisteringUsers
saveSyncIMUsersStatus(*status)
stats, err = srv_im.SyncDialogUsers()
status.Phase = phaseAddingFriends
status.LastStats = stats
saveSyncIMUsersStatus(*status)
return stats, err
}
func saveSyncIMUsersStatus(status SyncIMUsersTaskStatus) {
r := imSyncRedis()
if r == nil {
return
}
raw, err := json.Marshal(status)
if err != nil {
return
}
if err := r.Set(syncIMUsersStatusKey, string(raw), 0); err != nil {
log.Warn("saveSyncIMUsersStatus set failed", log.E(err))
}
}