refactor device service helpers

This commit is contained in:
ShukeBta
2026-06-26 14:25:27 +08:00
parent 09c4e9252c
commit 335a7fa14b
3 changed files with 355 additions and 333 deletions
+158
View File
@@ -0,0 +1,158 @@
package service
import (
"context"
"fmt"
"strings"
"time"
"go.uber.org/zap"
"github.com/ShukeBta/MediaStationGo/internal/model"
)
// SweepInactiveUsers is kept for compatibility with the existing scheduler; it
// now delegates to the custom account-cleanup policy.
func (s *DeviceService) SweepInactiveUsers(ctx context.Context) (int, error) {
return s.SweepAccountCleanup(ctx)
}
// SweepAccountCleanup runs the admin-defined account cleanup policy once.
// Users are kept when they satisfy any enabled keep rule. Users that do not
// meet any enabled rule are deleted.
func (s *DeviceService) SweepAccountCleanup(ctx context.Context) (int, error) {
cfg := loadBotConfig(ctx, s.repo)
if !cfg.AccountCleanupEnabled {
return 0, nil
}
candidates, err := s.accountCleanupCandidates(ctx, cfg)
if err != nil {
return 0, err
}
removed := 0
for _, candidate := range candidates {
s.notify(ctx, candidate.UserID, fmt.Sprintf("⛔️ 账号 <b>%s</b> 未满足保号规则,已被清理。\n规则结果:%s\n如需恢复请联系管理员。", candidate.Username, candidate.Details))
s.log.Warn("account cleanup: deleting account", zap.String("user", candidate.Username), zap.String("details", candidate.Details))
_ = s.repo.UserDevice.DeleteByUser(ctx, candidate.UserID)
if err := s.repo.User.Delete(ctx, candidate.UserID); err == nil {
removed++
}
}
return removed, nil
}
type accountCleanupCandidate struct {
UserID string
Username string
Details string
}
func (s *DeviceService) PreviewAccountCleanup(ctx context.Context) ([]accountCleanupCandidate, error) {
cfg := loadBotConfig(ctx, s.repo)
if !cfg.AccountCleanupEnabled {
return nil, nil
}
return s.accountCleanupCandidates(ctx, cfg)
}
func (s *DeviceService) accountCleanupCandidates(ctx context.Context, cfg botConfig) ([]accountCleanupCandidate, error) {
users, err := s.repo.User.List(ctx)
if err != nil {
return nil, err
}
candidates := make([]accountCleanupCandidate, 0)
for i := range users {
u := &users[i]
if s.isProtected(ctx, u) || !u.IsActive {
continue
}
keep, details := s.userMatchesCleanupPolicy(ctx, u, cfg)
if keep {
continue
}
candidates = append(candidates, accountCleanupCandidate{
UserID: u.ID,
Username: u.Username,
Details: details,
})
}
return candidates, nil
}
// randomWindowDays returns a random integer in [min,max]. The window is
// intentionally non-fixed per the operator's requirement ("随机触发").
func randomWindowDays(min, max int) int {
if min < 1 {
min = 1
}
if max < min {
max = min
}
if max == min {
return min
}
return min + secureRandomIntn(max-min+1)
}
func (s *DeviceService) userMatchesCleanupPolicy(ctx context.Context, u *model.User, cfg botConfig) (bool, string) {
rules := make([]accountCleanupRule, 0, len(cfg.AccountCleanupRules))
for _, r := range cfg.AccountCleanupRules {
if r.Enabled {
rules = append(rules, r)
}
}
if len(rules) == 0 {
return true, "无启用规则,跳过"
}
matches := 0
parts := make([]string, 0, len(rules))
for _, r := range rules {
ok, detail := s.userMatchesCleanupRule(ctx, u, r)
if ok {
matches++
parts = append(parts, "✅ "+detail)
} else {
parts = append(parts, "❌ "+detail)
}
}
required := 1
return matches >= required, fmt.Sprintf("满足 %d/%d 条,需要 %d 条;%s", matches, len(rules), required, strings.Join(parts, ";"))
}
func (s *DeviceService) userMatchesCleanupRule(ctx context.Context, u *model.User, r accountCleanupRule) (bool, string) {
switch r.Type {
case "watch_hours":
windowDays := randomWindowDays(r.WindowDaysMin, r.WindowDaysMax)
since := time.Now().Add(-time.Duration(windowDays) * 24 * time.Hour)
watched, _ := s.repo.UserDevice.WatchedMillisSince(ctx, u.ID, since)
hours := float64(watched) / 3600000
return hours >= r.MinHours, fmt.Sprintf("%s:近 %d 天观看 %.1f/%.1f 小时", r.Name, windowDays, hours, r.MinHours)
case "recent_login":
days := r.WindowDaysMax
if days < 1 {
days = r.WindowDaysMin
}
cutoff := time.Now().Add(-time.Duration(days) * 24 * time.Hour)
ok := u.LastLoginAt != nil && u.LastLoginAt.After(cutoff)
if !ok && s.sessions != nil {
ok = s.sessions.UserRecentlyActive(ctx, u.ID, time.Duration(days)*24*time.Hour)
}
return ok, fmt.Sprintf("%s:%d 天内登录", r.Name, days)
case "signin_streak":
rec, _ := s.repo.SignIn.Get(ctx, u.ID)
streak := 0
if rec != nil {
streak = rec.StreakDays
}
return streak >= r.MinCount, fmt.Sprintf("%s:连续签到 %d/%d 天", r.Name, streak, r.MinCount)
case "account_age_grace":
days := r.MinCount
if days < 1 {
days = r.WindowDaysMax
}
ok := u.CreatedAt.After(time.Now().Add(-time.Duration(days) * 24 * time.Hour))
return ok, fmt.Sprintf("%s:新账号 %d 天宽限", r.Name, days)
default:
return false, r.Name + ":未知规则"
}
}
+197
View File
@@ -0,0 +1,197 @@
package service
import (
"context"
"fmt"
"sort"
"strings"
"time"
"github.com/ShukeBta/MediaStationGo/internal/model"
)
// KickDevice marks a device as kicked so the next request from it is rejected
// (the client must log in again). Returns the affected device for messaging.
func (s *DeviceService) KickDevice(ctx context.Context, userID, deviceID string) error {
d, err := s.repo.UserDevice.Find(ctx, userID, deviceID)
if err != nil {
return err
}
if d == nil {
return fmt.Errorf("device not found")
}
if fp := strings.TrimSpace(d.Fingerprint); fp != "" {
return s.repo.UserDevice.SetKickedByFingerprint(ctx, userID, fp, true)
}
return s.repo.UserDevice.SetKicked(ctx, d.ID, true)
}
// KickAllDevices marks all devices for a user as kicked.
func (s *DeviceService) KickAllDevices(ctx context.Context, userID string) error {
return s.repo.UserDevice.SetKickedByUser(ctx, userID, true)
}
// ListDevices returns the device sessions for a user.
func (s *DeviceService) ListDevices(ctx context.Context, userID string) ([]model.UserDevice, error) {
rows, err := s.repo.UserDevice.ListByUser(ctx, userID)
if err != nil {
return nil, err
}
rows = collapseUserDeviceRows(rows)
if s.sessions == nil {
return rows, nil
}
now := s.sessions.now()
byDevice := make(map[string]int, len(rows))
byTerminal := make(map[string]int, len(rows))
for i := range rows {
byDevice[rows[i].DeviceID] = i
byTerminal[userDeviceTerminalKey(rows[i])] = i
}
for _, sess := range s.sessions.ListByUser(ctx, userID) {
online := sess.LastActivityAt.After(now.Add(-realtimeSessionOnlineTTL))
idx, ok := byDevice[sess.DeviceID]
if !ok {
idx, ok = byTerminal[sessionDeviceKey(sess)]
}
if ok {
if sess.LastActivityAt.After(rows[idx].LastSeenAt) {
rows[idx].LastSeenAt = sess.LastActivityAt
}
if sess.DeviceName != "" {
rows[idx].DeviceName = sess.DeviceName
}
if sess.Client != "" {
rows[idx].Client = sess.Client
}
if sess.RemoteEndPoint != "" {
rows[idx].LastIP = sess.RemoteEndPoint
}
if sess.LastPlaybackAt != nil {
rows[idx].LastPlayAt = sess.LastPlaybackAt
}
rows[idx].Realtime = true
rows[idx].Online = online
rows[idx].Playing = sess.IsPlaying && online
continue
}
row := model.UserDevice{
UserID: userID,
DeviceID: sess.DeviceID,
DeviceName: sess.DeviceName,
Client: sess.Client,
Fingerprint: fingerprint(sess.Client, sess.DeviceName),
LastIP: sess.RemoteEndPoint,
FirstSeenAt: sess.LastActivityAt,
LastSeenAt: sess.LastActivityAt,
LastPlayAt: sess.LastPlaybackAt,
Realtime: true,
Online: online,
Playing: sess.IsPlaying && online,
}
row.ID = "rt:" + sess.ID
rows = append(rows, row)
byDevice[row.DeviceID] = len(rows) - 1
byTerminal[userDeviceTerminalKey(row)] = len(rows) - 1
}
sort.SliceStable(rows, func(i, j int) bool {
return rows[i].LastSeenAt.After(rows[j].LastSeenAt)
})
return rows, nil
}
func collapseUserDeviceRows(rows []model.UserDevice) []model.UserDevice {
if len(rows) < 2 {
return rows
}
out := make([]model.UserDevice, 0, len(rows))
byTerminal := map[string]int{}
for _, row := range rows {
key := userDeviceTerminalKey(row)
if idx, ok := byTerminal[key]; ok {
out[idx] = mergeUserDeviceRows(out[idx], row)
continue
}
byTerminal[key] = len(out)
out = append(out, row)
}
return out
}
func mergeUserDeviceRows(a, b model.UserDevice) model.UserDevice {
if b.LastSeenAt.After(a.LastSeenAt) {
a.ID = b.ID
a.DeviceID = b.DeviceID
a.DeviceName = b.DeviceName
a.Client = b.Client
a.LastIP = b.LastIP
a.FirstSeenAt = earlierTime(a.FirstSeenAt, b.FirstSeenAt)
a.LastSeenAt = b.LastSeenAt
a.LastPlayAt = latestOptionalTime(a.LastPlayAt, b.LastPlayAt)
a.Fingerprint = firstNonEmptyString(b.Fingerprint, a.Fingerprint)
a.Kicked = b.Kicked
return a
}
a.FirstSeenAt = earlierTime(a.FirstSeenAt, b.FirstSeenAt)
a.LastPlayAt = latestOptionalTime(a.LastPlayAt, b.LastPlayAt)
a.Fingerprint = firstNonEmptyString(a.Fingerprint, b.Fingerprint)
return a
}
func userDeviceTerminalKey(row model.UserDevice) string {
if key := strings.TrimSpace(row.Fingerprint); key != "" {
return key
}
if strings.TrimSpace(row.DeviceName) != "" {
return "fp-" + fingerprint(row.Client, row.DeviceName)
}
if id := strings.TrimSpace(row.DeviceID); id != "" {
return id
}
return "unknown"
}
func earlierTime(a, b time.Time) time.Time {
if a.IsZero() || (!b.IsZero() && b.Before(a)) {
return b
}
return a
}
func latestOptionalTime(a, b *time.Time) *time.Time {
if a == nil {
return b
}
if b != nil && b.After(*a) {
return b
}
return a
}
// IsDeviceKicked reports whether a (user, device) pair was kicked and should be
// forced to re-authenticate.
func (s *DeviceService) IsDeviceKicked(ctx context.Context, userID, deviceID string) bool {
return s.IsTerminalKicked(ctx, userID, deviceID, "", "")
}
// IsTerminalKicked reports whether a request belongs to a kicked terminal.
func (s *DeviceService) IsTerminalKicked(ctx context.Context, userID, deviceID, deviceName, client string) bool {
if userID == "" {
return false
}
if strings.TrimSpace(deviceID) != "" {
d, err := s.repo.UserDevice.Find(ctx, userID, deviceID)
if err == nil && d != nil {
return d.Kicked
}
}
fp := ""
if strings.TrimSpace(deviceName) != "" {
fp = fingerprint(client, deviceName)
}
if fp == "" {
return false
}
d, err := s.repo.UserDevice.FindByFingerprint(ctx, userID, fp)
return err == nil && d != nil && d.Kicked
}
-333
View File
@@ -5,7 +5,6 @@ import (
"crypto/sha256"
"encoding/hex"
"fmt"
"sort"
"strings"
"time"
@@ -260,260 +259,6 @@ func (s *DeviceService) disableForPolicy(ctx context.Context, userID, reason str
s.log.Warn("device policy: disabled account", zap.String("user", u.Username), zap.String("reason", reason))
}
// SweepInactiveUsers is kept for compatibility with the existing scheduler; it
// now delegates to the custom account-cleanup policy.
func (s *DeviceService) SweepInactiveUsers(ctx context.Context) (int, error) {
return s.SweepAccountCleanup(ctx)
}
// SweepAccountCleanup runs the admin-defined account cleanup policy once.
// Users are kept when they satisfy any enabled keep rule. Users that do not
// meet any enabled rule are deleted.
func (s *DeviceService) SweepAccountCleanup(ctx context.Context) (int, error) {
cfg := loadBotConfig(ctx, s.repo)
if !cfg.AccountCleanupEnabled {
return 0, nil
}
candidates, err := s.accountCleanupCandidates(ctx, cfg)
if err != nil {
return 0, err
}
removed := 0
for _, candidate := range candidates {
s.notify(ctx, candidate.UserID, fmt.Sprintf("⛔️ 账号 <b>%s</b> 未满足保号规则,已被清理。\n规则结果:%s\n如需恢复请联系管理员。", candidate.Username, candidate.Details))
s.log.Warn("account cleanup: deleting account", zap.String("user", candidate.Username), zap.String("details", candidate.Details))
_ = s.repo.UserDevice.DeleteByUser(ctx, candidate.UserID)
if err := s.repo.User.Delete(ctx, candidate.UserID); err == nil {
removed++
}
}
return removed, nil
}
type accountCleanupCandidate struct {
UserID string
Username string
Details string
}
func (s *DeviceService) PreviewAccountCleanup(ctx context.Context) ([]accountCleanupCandidate, error) {
cfg := loadBotConfig(ctx, s.repo)
if !cfg.AccountCleanupEnabled {
return nil, nil
}
return s.accountCleanupCandidates(ctx, cfg)
}
func (s *DeviceService) accountCleanupCandidates(ctx context.Context, cfg botConfig) ([]accountCleanupCandidate, error) {
users, err := s.repo.User.List(ctx)
if err != nil {
return nil, err
}
candidates := make([]accountCleanupCandidate, 0)
for i := range users {
u := &users[i]
if s.isProtected(ctx, u) || !u.IsActive {
continue
}
keep, details := s.userMatchesCleanupPolicy(ctx, u, cfg)
if keep {
continue
}
candidates = append(candidates, accountCleanupCandidate{
UserID: u.ID,
Username: u.Username,
Details: details,
})
}
return candidates, nil
}
// KickDevice marks a device as kicked so the next request from it is rejected
// (the client must log in again). Returns the affected device for messaging.
func (s *DeviceService) KickDevice(ctx context.Context, userID, deviceID string) error {
d, err := s.repo.UserDevice.Find(ctx, userID, deviceID)
if err != nil {
return err
}
if d == nil {
return fmt.Errorf("device not found")
}
if fp := strings.TrimSpace(d.Fingerprint); fp != "" {
return s.repo.UserDevice.SetKickedByFingerprint(ctx, userID, fp, true)
}
return s.repo.UserDevice.SetKicked(ctx, d.ID, true)
}
// KickAllDevices marks all devices for a user as kicked.
func (s *DeviceService) KickAllDevices(ctx context.Context, userID string) error {
return s.repo.UserDevice.SetKickedByUser(ctx, userID, true)
}
// ListDevices returns the device sessions for a user.
func (s *DeviceService) ListDevices(ctx context.Context, userID string) ([]model.UserDevice, error) {
rows, err := s.repo.UserDevice.ListByUser(ctx, userID)
if err != nil {
return nil, err
}
rows = collapseUserDeviceRows(rows)
if s.sessions == nil {
return rows, nil
}
now := s.sessions.now()
byDevice := make(map[string]int, len(rows))
byTerminal := make(map[string]int, len(rows))
for i := range rows {
byDevice[rows[i].DeviceID] = i
byTerminal[userDeviceTerminalKey(rows[i])] = i
}
for _, sess := range s.sessions.ListByUser(ctx, userID) {
online := sess.LastActivityAt.After(now.Add(-realtimeSessionOnlineTTL))
idx, ok := byDevice[sess.DeviceID]
if !ok {
idx, ok = byTerminal[sessionDeviceKey(sess)]
}
if ok {
if sess.LastActivityAt.After(rows[idx].LastSeenAt) {
rows[idx].LastSeenAt = sess.LastActivityAt
}
if sess.DeviceName != "" {
rows[idx].DeviceName = sess.DeviceName
}
if sess.Client != "" {
rows[idx].Client = sess.Client
}
if sess.RemoteEndPoint != "" {
rows[idx].LastIP = sess.RemoteEndPoint
}
if sess.LastPlaybackAt != nil {
rows[idx].LastPlayAt = sess.LastPlaybackAt
}
rows[idx].Realtime = true
rows[idx].Online = online
rows[idx].Playing = sess.IsPlaying && online
continue
}
row := model.UserDevice{
UserID: userID,
DeviceID: sess.DeviceID,
DeviceName: sess.DeviceName,
Client: sess.Client,
Fingerprint: fingerprint(sess.Client, sess.DeviceName),
LastIP: sess.RemoteEndPoint,
FirstSeenAt: sess.LastActivityAt,
LastSeenAt: sess.LastActivityAt,
LastPlayAt: sess.LastPlaybackAt,
Realtime: true,
Online: online,
Playing: sess.IsPlaying && online,
}
row.ID = "rt:" + sess.ID
rows = append(rows, row)
byDevice[row.DeviceID] = len(rows) - 1
byTerminal[userDeviceTerminalKey(row)] = len(rows) - 1
}
sort.SliceStable(rows, func(i, j int) bool {
return rows[i].LastSeenAt.After(rows[j].LastSeenAt)
})
return rows, nil
}
func collapseUserDeviceRows(rows []model.UserDevice) []model.UserDevice {
if len(rows) < 2 {
return rows
}
out := make([]model.UserDevice, 0, len(rows))
byTerminal := map[string]int{}
for _, row := range rows {
key := userDeviceTerminalKey(row)
if idx, ok := byTerminal[key]; ok {
out[idx] = mergeUserDeviceRows(out[idx], row)
continue
}
byTerminal[key] = len(out)
out = append(out, row)
}
return out
}
func mergeUserDeviceRows(a, b model.UserDevice) model.UserDevice {
if b.LastSeenAt.After(a.LastSeenAt) {
a.ID = b.ID
a.DeviceID = b.DeviceID
a.DeviceName = b.DeviceName
a.Client = b.Client
a.LastIP = b.LastIP
a.FirstSeenAt = earlierTime(a.FirstSeenAt, b.FirstSeenAt)
a.LastSeenAt = b.LastSeenAt
a.LastPlayAt = latestOptionalTime(a.LastPlayAt, b.LastPlayAt)
a.Fingerprint = firstNonEmptyString(b.Fingerprint, a.Fingerprint)
a.Kicked = b.Kicked
return a
}
a.FirstSeenAt = earlierTime(a.FirstSeenAt, b.FirstSeenAt)
a.LastPlayAt = latestOptionalTime(a.LastPlayAt, b.LastPlayAt)
a.Fingerprint = firstNonEmptyString(a.Fingerprint, b.Fingerprint)
return a
}
func userDeviceTerminalKey(row model.UserDevice) string {
if key := strings.TrimSpace(row.Fingerprint); key != "" {
return key
}
if strings.TrimSpace(row.DeviceName) != "" {
return "fp-" + fingerprint(row.Client, row.DeviceName)
}
if id := strings.TrimSpace(row.DeviceID); id != "" {
return id
}
return "unknown"
}
func earlierTime(a, b time.Time) time.Time {
if a.IsZero() || (!b.IsZero() && b.Before(a)) {
return b
}
return a
}
func latestOptionalTime(a, b *time.Time) *time.Time {
if a == nil {
return b
}
if b != nil && b.After(*a) {
return b
}
return a
}
// IsDeviceKicked reports whether a (user, device) pair was kicked and should be
// forced to re-authenticate.
func (s *DeviceService) IsDeviceKicked(ctx context.Context, userID, deviceID string) bool {
return s.IsTerminalKicked(ctx, userID, deviceID, "", "")
}
// IsTerminalKicked reports whether a request belongs to a kicked terminal.
func (s *DeviceService) IsTerminalKicked(ctx context.Context, userID, deviceID, deviceName, client string) bool {
if userID == "" {
return false
}
if strings.TrimSpace(deviceID) != "" {
d, err := s.repo.UserDevice.Find(ctx, userID, deviceID)
if err == nil && d != nil {
return d.Kicked
}
}
fp := ""
if strings.TrimSpace(deviceName) != "" {
fp = fingerprint(client, deviceName)
}
if fp == "" {
return false
}
d, err := s.repo.UserDevice.FindByFingerprint(ctx, userID, fp)
return err == nil && d != nil && d.Kicked
}
func (s *DeviceService) UserRecentlyActive(ctx context.Context, userID string, within time.Duration) bool {
return s.sessions != nil && s.sessions.UserRecentlyActive(ctx, userID, within)
}
@@ -531,84 +276,6 @@ func (s *DeviceService) notify(ctx context.Context, userID, text string) {
}
}
// randomWindowDays returns a random integer in [min,max]. The window is
// intentionally non-fixed per the operator's requirement ("随机触发").
func randomWindowDays(min, max int) int {
if min < 1 {
min = 1
}
if max < min {
max = min
}
if max == min {
return min
}
return min + secureRandomIntn(max-min+1)
}
func (s *DeviceService) userMatchesCleanupPolicy(ctx context.Context, u *model.User, cfg botConfig) (bool, string) {
rules := make([]accountCleanupRule, 0, len(cfg.AccountCleanupRules))
for _, r := range cfg.AccountCleanupRules {
if r.Enabled {
rules = append(rules, r)
}
}
if len(rules) == 0 {
return true, "无启用规则,跳过"
}
matches := 0
parts := make([]string, 0, len(rules))
for _, r := range rules {
ok, detail := s.userMatchesCleanupRule(ctx, u, r)
if ok {
matches++
parts = append(parts, "✅ "+detail)
} else {
parts = append(parts, "❌ "+detail)
}
}
required := 1
return matches >= required, fmt.Sprintf("满足 %d/%d 条,需要 %d 条;%s", matches, len(rules), required, strings.Join(parts, ";"))
}
func (s *DeviceService) userMatchesCleanupRule(ctx context.Context, u *model.User, r accountCleanupRule) (bool, string) {
switch r.Type {
case "watch_hours":
windowDays := randomWindowDays(r.WindowDaysMin, r.WindowDaysMax)
since := time.Now().Add(-time.Duration(windowDays) * 24 * time.Hour)
watched, _ := s.repo.UserDevice.WatchedMillisSince(ctx, u.ID, since)
hours := float64(watched) / 3600000
return hours >= r.MinHours, fmt.Sprintf("%s:近 %d 天观看 %.1f/%.1f 小时", r.Name, windowDays, hours, r.MinHours)
case "recent_login":
days := r.WindowDaysMax
if days < 1 {
days = r.WindowDaysMin
}
cutoff := time.Now().Add(-time.Duration(days) * 24 * time.Hour)
ok := u.LastLoginAt != nil && u.LastLoginAt.After(cutoff)
if !ok && s.sessions != nil {
ok = s.sessions.UserRecentlyActive(ctx, u.ID, time.Duration(days)*24*time.Hour)
}
return ok, fmt.Sprintf("%s:%d 天内登录", r.Name, days)
case "signin_streak":
rec, _ := s.repo.SignIn.Get(ctx, u.ID)
streak := 0
if rec != nil {
streak = rec.StreakDays
}
return streak >= r.MinCount, fmt.Sprintf("%s:连续签到 %d/%d 天", r.Name, streak, r.MinCount)
case "account_age_grace":
days := r.MinCount
if days < 1 {
days = r.WindowDaysMax
}
ok := u.CreatedAt.After(time.Now().Add(-time.Duration(days) * 24 * time.Hour))
return ok, fmt.Sprintf("%s:新账号 %d 天宽限", r.Name, days)
default:
return false, r.Name + ":未知规则"
}
}
func deviceLabel(name, client string) string {
name = strings.TrimSpace(name)
client = strings.TrimSpace(client)