feat(bot): manage device policy via telegram commands

This commit is contained in:
ShukeBta
2026-06-07 12:24:21 +08:00
parent cf2f77f35d
commit b38d47a27d
11 changed files with 667 additions and 237 deletions
+24
View File
@@ -698,6 +698,30 @@ func (r *RefreshTokenRepository) RevokeByUserID(ctx context.Context, userID stri
Where("user_id = ?", userID).Update("revoked", true).Error
}
// RevokeOldestActiveByUserID keeps at most limit active refresh tokens for a
// user by revoking the oldest non-expired, non-revoked tokens.
func (r *RefreshTokenRepository) RevokeOldestActiveByUserID(ctx context.Context, userID string, limit int) error {
if limit < 1 {
limit = 1
}
var tokens []model.RefreshToken
if err := r.db.WithContext(ctx).
Where("user_id = ? AND revoked = ? AND expires_at > ?", userID, false, time.Now()).
Order("created_at desc, id desc").
Find(&tokens).Error; err != nil {
return err
}
if len(tokens) <= limit {
return nil
}
ids := make([]string, 0, len(tokens)-limit)
for _, token := range tokens[limit:] {
ids = append(ids, token.ID)
}
return r.db.WithContext(ctx).Model(&model.RefreshToken{}).
Where("id IN ?", ids).Update("revoked", true).Error
}
// DeleteExpired removes all expired refresh tokens.
func (r *RefreshTokenRepository) DeleteExpired(ctx context.Context) error {
return r.db.WithContext(ctx).Where("expires_at < ?", time.Now()).Delete(&model.RefreshToken{}).Error
+28
View File
@@ -188,6 +188,34 @@ func TestAdminResetPasswordAllowsLoginWithNewPassword(t *testing.T) {
}
}
func TestLoginKeepsOnlyConfiguredActiveRefreshTokens(t *testing.T) {
ctx := context.Background()
repos, auth, _, _ := newAuthTestServices(t)
user, _, err := auth.Register(ctx, "viewer", "password")
if err != nil {
t.Fatalf("register: %v", err)
}
if err := repos.Setting.Set(ctx, SettingMaxLoggedClients, "3"); err != nil {
t.Fatalf("set max clients: %v", err)
}
for i := 0; i < 5; i++ {
if _, err := auth.Login(ctx, "viewer", "password"); err != nil {
t.Fatalf("login %d: %v", i+1, err)
}
}
var active int64
if err := repos.DB.Model(&model.RefreshToken{}).
Where("user_id = ? AND revoked = ? AND expires_at > ?", user.ID, false, time.Now()).
Count(&active).Error; err != nil {
t.Fatalf("count active refresh tokens: %v", err)
}
if active != 3 {
t.Fatalf("active refresh tokens should be capped at 3, got %d", active)
}
}
func TestDefaultPermissionsAreViewerOnly(t *testing.T) {
perms := DefaultPermissions("user-1")
if !perms.CanViewDashboard || !perms.CanPlayMedia || !perms.CanExternalPlayer {
+111
View File
@@ -2,6 +2,7 @@ package service
import (
"context"
"strings"
"testing"
"time"
@@ -262,3 +263,113 @@ func TestProtectedAdminNeverViolated(t *testing.T) {
t.Fatalf("admin should accrue no warnings, got %d", got.ShareWarnings)
}
}
func TestBotAdminCommandsManageDevicePolicy(t *testing.T) {
ctx := context.Background()
repos, bot := newBotTestService(t)
admin := &model.User{Username: "root", PasswordHash: "x", Role: "admin", IsActive: true}
if err := repos.User.Create(ctx, admin); err != nil {
t.Fatal(err)
}
if err := repos.DB.Create(&model.TelegramBinding{
TelegramUserID: 9001,
TelegramName: "@root",
ChatID: 9001,
UserID: admin.ID,
}).Error; err != nil {
t.Fatal(err)
}
channel := &model.NotifyChannel{Name: "Telegram", Type: "telegram", Enabled: true, Config: `{"admin_user_ids":"9001"}`}
msg := &TelegramMessage{From: TelegramUser{ID: 9001, Username: "root"}, Chat: TelegramChat{ID: 9001, Type: "private"}}
reply, err := bot.executeCommand(ctx, channel, msg, "/antishare on play=4 login=5 warn=3")
if err != nil {
t.Fatal(err)
}
if !strings.Contains(reply.Text, "防共享:<b>已开启</b>") {
t.Fatalf("expected antishare enabled reply, got %q", reply.Text)
}
cfg := loadBotConfig(ctx, repos)
if !cfg.AntiShareEnabled || cfg.MaxConcurrentPlay != 4 || cfg.MaxLoggedClients != 5 || cfg.WarnThreshold != 3 {
t.Fatalf("unexpected device policy: %+v", cfg)
}
reply, err = bot.executeCommand(ctx, channel, msg, "/cleanup_mode count 2")
if err != nil {
t.Fatal(err)
}
cfg = loadBotConfig(ctx, repos)
if cfg.AccountCleanupKeepMode != "count" || cfg.AccountCleanupRequiredCount != 2 {
t.Fatalf("unexpected cleanup mode: %+v; reply=%q", cfg, reply.Text)
}
reply, err = bot.executeCommand(ctx, channel, msg, "/cleanup_rule add recent_login login_7d 七天内登录 7")
if err != nil {
t.Fatal(err)
}
cfg = loadBotConfig(ctx, repos)
found := false
for _, rule := range cfg.AccountCleanupRules {
if rule.ID == "login_7d" && rule.Type == "recent_login" && rule.WindowDaysMax == 7 {
found = true
}
}
if !found {
t.Fatalf("cleanup rule not added; reply=%q rules=%+v", reply.Text, cfg.AccountCleanupRules)
}
}
func TestBotUserCommandsAndAdminGate(t *testing.T) {
ctx := context.Background()
repos, bot := newBotTestService(t)
user := &model.User{Username: "viewer", PasswordHash: "x", Role: "user", IsActive: true}
if err := repos.User.Create(ctx, user); err != nil {
t.Fatal(err)
}
if err := repos.DB.Create(&model.TelegramBinding{
TelegramUserID: 9101,
TelegramName: "@viewer",
ChatID: 9101,
UserID: user.ID,
}).Error; err != nil {
t.Fatal(err)
}
now := time.Now()
if err := repos.UserDevice.Create(ctx, &model.UserDevice{
UserID: user.ID, DeviceID: "dev-1", DeviceName: "iPhone", Client: "Infuse", FirstSeenAt: now, LastSeenAt: now,
}); err != nil {
t.Fatal(err)
}
channel := &model.NotifyChannel{Name: "Telegram", Type: "telegram", Enabled: true, Config: `{"admin_user_ids":"9001"}`}
msg := &TelegramMessage{From: TelegramUser{ID: 9101, Username: "viewer"}, Chat: TelegramChat{ID: 9101, Type: "private"}}
reply, err := bot.executeCommand(ctx, channel, msg, "/antishare on")
if err != nil {
t.Fatal(err)
}
if !strings.Contains(reply.Text, "仅管理员") {
t.Fatalf("regular user should not manage policy, got %q", reply.Text)
}
reply, err = bot.executeCommand(ctx, channel, msg, "/devices")
if err != nil {
t.Fatal(err)
}
if !strings.Contains(reply.Text, "我的登录设备") {
t.Fatalf("expected device list, got %q", reply.Text)
}
reply, err = bot.executeCommand(ctx, channel, msg, "/kick 1")
if err != nil {
t.Fatal(err)
}
if !strings.Contains(reply.Text, "已踢下线") {
t.Fatalf("expected kick feedback, got %q", reply.Text)
}
if kicked := bot.device; kicked != nil {
t.Fatal("test should not require wired device service")
}
if ok := NewDeviceService(zap.NewNop(), repos).IsDeviceKicked(ctx, user.ID, "dev-1"); !ok {
t.Fatal("device should be marked kicked")
}
}
+2 -2
View File
@@ -9,8 +9,8 @@ import (
"github.com/ShukeBta/MediaStationGo/internal/repository"
)
// Bot / 设备管控相关的设置键。全部存储在 settings 表,可由管理员在 Bot 或
// 系统设置页调整。带安全默认值:所有"自动删号"策略默认关闭。
// Bot / 设备管控相关的设置键。全部存储在 settings 表,由管理员通过
// Telegram Bot 命令调整。带安全默认值:所有"自动删号"策略默认关闭。
const (
// 开放注册(开注名额)。
SettingOpenRegEnabled = "telegram.openreg_enabled" // 是否开放注册
+1 -1
View File
@@ -239,7 +239,7 @@ func (c *Container) Boot() {
// 启动调度器定时任务
c.Scheduler.Start(c.stopCtx)
// 账号删号/保号规则巡检:默认关闭,由管理员在 Bot 或运维工具开启。
// 账号删号/保号规则巡检:默认关闭,由管理员通过 Telegram Bot 命令开启。
// 每天触发一次评估;规则里的窗口可随机,不固定。
if c.Device != nil {
go c.runInactivitySweeper(c.stopCtx)
+90 -3
View File
@@ -134,6 +134,8 @@ func NewTelegramBotService(log *zap.Logger, repo *repository.Container, crypto *
// 默认关闭,只有管理员在系统设置 / Bot 管理命令中显式开启后才允许注册。
const TelegramRegistrationSettingKey = "telegram.registration_enabled"
var errTelegramAccountAlreadyBound = errors.New("该媒体账号已绑定其他 Telegram,请联系管理员解绑")
// registrationEnabled 读取注册开关;默认关闭。
func (s *TelegramBotService) registrationEnabled(ctx context.Context) bool {
v, err := s.repo.Setting.Get(ctx, TelegramRegistrationSettingKey)
@@ -251,13 +253,58 @@ func (s *TelegramBotService) executeCommand(ctx context.Context, channel *model.
return telegramCommandReply{Text: s.cmdHelp(ctx, msg)}, nil
case "/hideadult", "/hide_adult", "/adult":
return s.cmdHideAdult(ctx, msg, args), nil
case "/account", "/me":
return s.replyAccount(ctx, msg), nil
case "/devices":
return s.replyDevices(ctx, msg), nil
case "/kick":
return s.cmdKick(ctx, msg, args), nil
case "/setname", "/rename":
return s.cmdSetName(ctx, msg, args), nil
case "/setpass", "/passwd", "/password":
return s.cmdSetPass(ctx, msg, args), nil
case "/register", "/reg", "/signup":
return s.cmdRegister(ctx, channel, msg, args), nil
case "/registration", "/reg_switch":
case "/registration", "/reg_switch", "/openreg":
if !s.telegramUserIsAdmin(ctx, channel, msg.From.ID) {
return telegramCommandReply{Text: "此命令仅管理员可用。"}, nil
}
return s.cmdRegistrationToggle(ctx, args), nil
case "/devicepolicy", "/policy":
if !s.telegramUserIsAdmin(ctx, channel, msg.From.ID) {
return telegramCommandReply{Text: "此命令仅管理员可用。"}, nil
}
return s.cmdDevicePolicy(ctx, args), nil
case "/antishare":
if !s.telegramUserIsAdmin(ctx, channel, msg.From.ID) {
return telegramCommandReply{Text: "此命令仅管理员可用。"}, nil
}
return s.cmdAntiShare(ctx, args), nil
case "/cleanup":
if !s.telegramUserIsAdmin(ctx, channel, msg.From.ID) {
return telegramCommandReply{Text: "此命令仅管理员可用。"}, nil
}
return s.cmdCleanup(ctx, args), nil
case "/cleanup_mode":
if !s.telegramUserIsAdmin(ctx, channel, msg.From.ID) {
return telegramCommandReply{Text: "此命令仅管理员可用。"}, nil
}
return s.cmdCleanupMode(ctx, args), nil
case "/cleanup_rule":
if !s.telegramUserIsAdmin(ctx, channel, msg.From.ID) {
return telegramCommandReply{Text: "此命令仅管理员可用。"}, nil
}
return s.cmdCleanupRule(ctx, args), nil
case "/ban":
if !s.telegramUserIsAdmin(ctx, channel, msg.From.ID) {
return telegramCommandReply{Text: "此命令仅管理员可用。"}, nil
}
return s.cmdUserBan(ctx, args, false), nil
case "/unban":
if !s.telegramUserIsAdmin(ctx, channel, msg.From.ID) {
return telegramCommandReply{Text: "此命令仅管理员可用。"}, nil
}
return s.cmdUserBan(ctx, args, true), nil
case "/status":
if !s.telegramUserIsAdmin(ctx, channel, msg.From.ID) {
return telegramCommandReply{Text: "此命令仅管理员可用。普通用户只能使用 /start 绑定账号,并通过按钮隐藏成人目录。"}, nil
@@ -305,8 +352,10 @@ func telegramCommandName(text string) string {
func telegramSupportedCommand(cmd string) bool {
switch cmd {
case "/start", "/menu", "/cancel", "/help", "/hideadult", "/hide_adult", "/adult",
"/register", "/reg", "/signup", "/registration", "/reg_switch",
"/status", "/search", "/downloads", "/stats":
"/account", "/me", "/devices", "/kick", "/setname", "/rename", "/setpass", "/passwd", "/password",
"/register", "/reg", "/signup", "/registration", "/reg_switch", "/openreg",
"/devicepolicy", "/policy", "/antishare", "/cleanup", "/cleanup_mode", "/cleanup_rule",
"/ban", "/unban", "/status", "/search", "/downloads", "/stats":
return true
default:
return false
@@ -461,14 +510,26 @@ func (s *TelegramBotService) cmdHelp(ctx context.Context, msg *TelegramMessage)
return "<b>MediaStationGo 用户命令</b>\n\n" +
register +
"<b>/start 用户名 密码</b> — 绑定账号\n" +
"<b>/account</b> — 查看账号状态\n" +
"<b>/devices</b> — 查看登录设备\n" +
"<b>/kick all|编号</b> — 踢下线设备\n" +
"<b>/setname 新用户名</b> — 修改用户名\n" +
"<b>/setpass 新密码</b> — 修改密码\n" +
"<b>/hideadult on|off</b> — 隐藏或显示成人目录\n\n" +
"系统状态、搜索、下载列表与统计命令仅管理员可用。"
}
return "<b>MediaStationGo 命令列表</b>\n\n" +
"<b>/start</b> — 开始使用\n" +
"<b>/help</b> — 帮助信息\n" +
"<b>/account</b> / <b>/devices</b> / <b>/kick all|编号</b> — 用户自助设备管理\n" +
"<b>/setname 新用户名</b> / <b>/setpass 新密码</b> — 用户自助改名改密\n" +
"<b>/register 用户名 密码</b> — 注册新账号(需管理员开启)\n" +
"<b>/registration on|off</b> — 开启/关闭普通用户注册(管理员)\n" +
"<b>/antishare on play=3 login=3 warn=2</b> — 防共享策略(管理员)\n" +
"<b>/cleanup on|off|run</b> — 删号规则开关/巡检(管理员)\n" +
"<b>/cleanup_mode any|all|count 2</b> — 保号模式(管理员)\n" +
"<b>/cleanup_rule list|add|del|enable|disable</b> — 保号规则(管理员)\n" +
"<b>/ban 用户名</b> / <b>/unban 用户名</b> — 禁用/解禁用户(管理员)\n" +
"<b>/hideadult on|off</b> — 隐藏/显示当前绑定账号的成人目录\n" +
"<b>/status</b> — 系统运行状态\n" +
"<b>/search 关键词</b> — 搜索媒体库\n" +
@@ -1102,6 +1163,11 @@ func (s *TelegramBotService) upsertTelegramBinding(ctx context.Context, msg *Tel
var existing model.TelegramBinding
err := s.repo.DB.WithContext(ctx).Where("telegram_user_id = ?", int64(msg.From.ID)).First(&existing).Error
if err == nil {
if existing.UserID != userID {
if err := s.ensureTelegramAccountBindingAvailable(ctx, userID, int64(msg.From.ID)); err != nil {
return err
}
}
return s.repo.DB.WithContext(ctx).Model(&existing).Updates(map[string]any{
"telegram_name": name,
"chat_id": int64(msg.Chat.ID),
@@ -1114,6 +1180,9 @@ func (s *TelegramBotService) upsertTelegramBinding(ctx context.Context, msg *Tel
if err := s.repo.DB.WithContext(ctx).Unscoped().Where("telegram_user_id = ?", int64(msg.From.ID)).Delete(&model.TelegramBinding{}).Error; err != nil {
return err
}
if err := s.ensureTelegramAccountBindingAvailable(ctx, userID, int64(msg.From.ID)); err != nil {
return err
}
return s.repo.DB.WithContext(ctx).Create(&model.TelegramBinding{
TelegramUserID: int64(msg.From.ID),
TelegramName: name,
@@ -1122,6 +1191,24 @@ func (s *TelegramBotService) upsertTelegramBinding(ctx context.Context, msg *Tel
}).Error
}
func (s *TelegramBotService) ensureTelegramAccountBindingAvailable(ctx context.Context, userID string, telegramUserID int64) error {
var bound model.TelegramBinding
err := s.repo.DB.WithContext(ctx).
Where("user_id = ? AND telegram_user_id <> ?", userID, telegramUserID).
First(&bound).Error
if errors.Is(err, gorm.ErrRecordNotFound) {
return nil
}
if err != nil {
return err
}
if user, _ := s.repo.User.FindByID(ctx, bound.UserID); user == nil {
_ = s.repo.DB.WithContext(ctx).Unscoped().Delete(&model.TelegramBinding{}, "id = ?", bound.ID).Error
return nil
}
return errTelegramAccountAlreadyBound
}
func parseStartCredentials(args []string) (string, string) {
if len(args) >= 2 {
return strings.TrimSpace(args[0]), strings.TrimSpace(strings.Join(args[1:], " "))
@@ -161,3 +161,48 @@ func TestTelegramStartClearsStaleUserBinding(t *testing.T) {
t.Fatalf("stale binding should be removed, got %d", count)
}
}
func TestTelegramStartRejectsAccountAlreadyBoundToAnotherTelegram(t *testing.T) {
ctx := t.Context()
repos, auth, _, _ := newAuthTestServices(t)
user, _, err := auth.Register(ctx, "viewer", "secret-pass")
if err != nil {
t.Fatalf("register: %v", err)
}
if err := repos.DB.Create(&model.TelegramBinding{
TelegramUserID: 20001,
TelegramName: "@viewer-one",
ChatID: 20001,
UserID: user.ID,
}).Error; err != nil {
t.Fatalf("create binding: %v", err)
}
bot := NewTelegramBotService(zap.NewNop(), repos, nil, auth)
cfgJSON, _ := json.Marshal(map[string]string{"admin_user_ids": "20002"})
if err := repos.DB.AutoMigrate(&model.NotifyChannel{}); err != nil {
t.Fatalf("migrate notify channel: %v", err)
}
if err := repos.DB.Create(&model.NotifyChannel{Name: "Telegram", Type: "telegram", Enabled: true, Config: string(cfgJSON)}).Error; err != nil {
t.Fatalf("create notify channel: %v", err)
}
msg := &TelegramMessage{
From: TelegramUser{ID: 20002, Username: "viewer-two", FirstName: "Viewer Two"},
Chat: TelegramChat{ID: 20002, Type: "private"},
}
reply := bot.cmdStart(ctx, msg, []string{"viewer", "secret-pass"})
if !strings.Contains(reply.Text, "已绑定其他 Telegram") {
t.Fatalf("expected already-bound rejection, got %q", reply.Text)
}
var accountBindings int64
if err := repos.DB.Model(&model.TelegramBinding{}).Where("user_id = ?", user.ID).Count(&accountBindings).Error; err != nil {
t.Fatalf("count account bindings: %v", err)
}
if accountBindings != 1 {
t.Fatalf("account should keep exactly one telegram binding, got %d", accountBindings)
}
if binding := bot.telegramBinding(ctx, 20002); binding != nil {
t.Fatal("second telegram account must not be bound")
}
}
+354 -1
View File
@@ -2,6 +2,7 @@ package service
import (
"context"
"encoding/json"
"fmt"
"strconv"
"strings"
@@ -216,6 +217,63 @@ func (s *TelegramBotService) handlePendingText(ctx context.Context, channel *mod
// ── 用户自助 ──────────────────────────────────────────────────────────────
func (s *TelegramBotService) cmdKick(ctx context.Context, msg *TelegramMessage, args []string) telegramCommandReply {
user := s.boundUser(ctx, msg.From.ID)
if user == nil {
return telegramCommandReply{Text: "请先绑定账号:<code>/start 用户名 密码</code>"}
}
if len(args) == 0 {
return telegramCommandReply{Text: "请指定要踢下线的设备:<code>/kick all</code> 或 <code>/kick 设备编号</code>。先用 <code>/devices</code> 查看编号。"}
}
target := strings.TrimSpace(args[0])
if strings.EqualFold(target, "all") || target == "全部" {
if s.device != nil {
if err := s.device.KickAllDevices(ctx, user.ID); err != nil {
return telegramCommandReply{Text: "踢下线失败:" + err.Error()}
}
} else if err := s.repo.UserDevice.SetKickedByUser(ctx, user.ID, true); err != nil {
return telegramCommandReply{Text: "踢下线失败:" + err.Error()}
}
return telegramCommandReply{Text: "已踢下线此账号的全部设备。"}
}
devices, _ := s.repo.UserDevice.ListByUser(ctx, user.ID)
if len(devices) == 0 {
return telegramCommandReply{Text: "当前没有记录到登录设备。"}
}
var chosen *model.UserDevice
if n, err := strconv.Atoi(target); err == nil && n >= 1 && n <= len(devices) {
chosen = &devices[n-1]
} else {
for i := range devices {
if devices[i].ID == target || devices[i].DeviceID == target {
chosen = &devices[i]
break
}
}
}
if chosen == nil {
return telegramCommandReply{Text: "未找到该设备。请用 <code>/devices</code> 查看设备编号后重试。"}
}
if err := s.repo.UserDevice.SetKicked(ctx, chosen.ID, true); err != nil {
return telegramCommandReply{Text: "踢下线失败:" + err.Error()}
}
return telegramCommandReply{Text: fmt.Sprintf("已踢下线:<b>%s</b>。", deviceLabel(chosen.DeviceName, chosen.Client))}
}
func (s *TelegramBotService) cmdSetName(ctx context.Context, msg *TelegramMessage, args []string) telegramCommandReply {
if len(args) == 0 {
return telegramCommandReply{Text: "请发送:<code>/setname 新用户名</code>"}
}
return s.selfSetName(ctx, msg, strings.Join(args, " "))
}
func (s *TelegramBotService) cmdSetPass(ctx context.Context, msg *TelegramMessage, args []string) telegramCommandReply {
if len(args) == 0 {
return telegramCommandReply{Text: "请发送:<code>/setpass 新密码</code>"}
}
return s.selfSetPass(ctx, msg, strings.Join(args, " "))
}
func (s *TelegramBotService) replyAccount(ctx context.Context, msg *TelegramMessage) telegramCommandReply {
user := s.boundUser(ctx, msg.From.ID)
if user == nil {
@@ -583,7 +641,7 @@ func (s *TelegramBotService) protectReason(ctx context.Context, userID string) s
func (s *TelegramBotService) replyDevicePolicy(ctx context.Context) telegramCommandReply {
cfg := loadBotConfig(ctx, s.repo)
text := fmt.Sprintf(
"<b>设备策略</b>\n\n① 防共享:<b>%s</b>\n 并发播放上限 %d / 登录客户端上限 %d;超限会禁用账号,管理员可解禁。\n 设备指纹异常警告 %d 次后禁用账号。\n\n② 自定义删号规则:<b>%s</b>\n 保号模式:%s;需要满足 %d 条;启用规则 %d 条。\n\n策略默认关闭;删号前会先通过 Bot 通知用户;管理员/受保护账号永不自动处理。",
"<b>设备策略</b>\n\n① 防共享:<b>%s</b>\n 并发播放上限 %d / 登录客户端上限 %d;超限会禁用账号,管理员可解禁。\n 设备指纹异常警告 %d 次后禁用账号。\n\n② 自定义删号规则:<b>%s</b>\n 保号模式:%s;需要满足 %d 条;启用规则 %d 条。\n\n<b>命令:</b>\n<code>/antishare on play=3 login=3 warn=2</code>\n<code>/cleanup on|off|run</code>\n<code>/cleanup_mode any|all|count 2</code>\n<code>/cleanup_rule list|add|del|enable|disable</code>\n\n策略默认关闭;删号前会先通过 Bot 通知用户;管理员/受保护账号永不自动处理。",
onOff(cfg.AntiShareEnabled), cfg.MaxConcurrentPlay, cfg.MaxLoggedClients, cfg.WarnThreshold,
onOff(cfg.AccountCleanupEnabled), cleanupModeLabel(cfg.AccountCleanupKeepMode), cfg.AccountCleanupRequiredCount, countEnabledCleanupRules(cfg.AccountCleanupRules))
return telegramCommandReply{
@@ -596,6 +654,168 @@ func (s *TelegramBotService) replyDevicePolicy(ctx context.Context) telegramComm
}
}
func (s *TelegramBotService) cmdDevicePolicy(ctx context.Context, args []string) telegramCommandReply {
if len(args) == 0 || strings.EqualFold(args[0], "status") {
return s.replyDevicePolicy(ctx)
}
switch strings.ToLower(strings.TrimSpace(args[0])) {
case "run", "sweep":
return s.cmdCleanup(ctx, []string{"run"})
default:
return telegramCommandReply{Text: "用法:<code>/devicepolicy</code> 查看策略,或使用 <code>/antishare</code>、<code>/cleanup</code>、<code>/cleanup_rule</code> 管理。"}
}
}
func (s *TelegramBotService) cmdAntiShare(ctx context.Context, args []string) telegramCommandReply {
if len(args) == 0 || strings.EqualFold(args[0], "status") {
return s.replyDevicePolicy(ctx)
}
enabled, ok := parseCommandBool(args[0])
if !ok {
return telegramCommandReply{Text: "用法:<code>/antishare on|off [play=3] [login=3] [warn=2]</code>"}
}
if err := s.repo.Setting.Set(ctx, SettingAntiShareEnabled, strconv.FormatBool(enabled)); err != nil {
return telegramCommandReply{Text: "更新失败:" + err.Error()}
}
for _, arg := range args[1:] {
key, value, ok := strings.Cut(arg, "=")
if !ok {
continue
}
n, err := strconv.Atoi(strings.TrimSpace(value))
if err != nil || n < 1 {
continue
}
switch strings.ToLower(strings.TrimSpace(key)) {
case "play", "maxplay", "播放":
_ = s.repo.Setting.Set(ctx, SettingMaxConcurrentPlay, strconv.Itoa(n))
case "login", "client", "clients", "登录":
_ = s.repo.Setting.Set(ctx, SettingMaxLoggedClients, strconv.Itoa(n))
case "warn", "warnings", "警告":
_ = s.repo.Setting.Set(ctx, SettingWarnThreshold, strconv.Itoa(n))
}
}
return s.replyDevicePolicy(ctx)
}
func (s *TelegramBotService) cmdCleanup(ctx context.Context, args []string) telegramCommandReply {
if len(args) == 0 || strings.EqualFold(args[0], "status") {
return s.replyDevicePolicy(ctx)
}
switch strings.ToLower(strings.TrimSpace(args[0])) {
case "on", "true", "1", "开启", "enable":
if err := s.repo.Setting.Set(ctx, SettingAccountCleanupEnabled, "true"); err != nil {
return telegramCommandReply{Text: "开启失败:" + err.Error()}
}
return s.replyDevicePolicy(ctx)
case "off", "false", "0", "关闭", "disable":
if err := s.repo.Setting.Set(ctx, SettingAccountCleanupEnabled, "false"); err != nil {
return telegramCommandReply{Text: "关闭失败:" + err.Error()}
}
return s.replyDevicePolicy(ctx)
case "run", "sweep", "巡检":
device := s.device
if device == nil {
device = NewDeviceService(s.log, s.repo)
}
removed, err := device.SweepAccountCleanup(ctx)
if err != nil {
return telegramCommandReply{Text: "巡检失败:" + err.Error()}
}
return telegramCommandReply{Text: fmt.Sprintf("删号规则巡检完成,清理 <b>%d</b> 个账号。", removed)}
default:
return telegramCommandReply{Text: "用法:<code>/cleanup on|off|run</code>"}
}
}
func (s *TelegramBotService) cmdCleanupMode(ctx context.Context, args []string) telegramCommandReply {
if len(args) == 0 {
return telegramCommandReply{Text: "用法:<code>/cleanup_mode any</code>、<code>/cleanup_mode all</code> 或 <code>/cleanup_mode count 2</code>"}
}
mode := strings.ToLower(strings.TrimSpace(args[0]))
if mode != "any" && mode != "all" && mode != "count" {
return telegramCommandReply{Text: "保号模式无效,只支持 any / all / count。"}
}
if err := s.repo.Setting.Set(ctx, SettingAccountCleanupKeepMode, mode); err != nil {
return telegramCommandReply{Text: "更新失败:" + err.Error()}
}
if mode == "count" && len(args) > 1 {
n, err := strconv.Atoi(args[1])
if err == nil && n > 0 {
_ = s.repo.Setting.Set(ctx, SettingAccountCleanupRequiredCount, strconv.Itoa(n))
}
}
return s.replyDevicePolicy(ctx)
}
func (s *TelegramBotService) cmdCleanupRule(ctx context.Context, args []string) telegramCommandReply {
if len(args) == 0 {
return telegramCommandReply{Text: cleanupRuleHelp()}
}
rules := s.currentCleanupRules(ctx)
action := strings.ToLower(strings.TrimSpace(args[0]))
switch action {
case "list", "ls", "status":
return telegramCommandReply{Text: formatCleanupRules(rules)}
case "del", "delete", "rm":
if len(args) < 2 {
return telegramCommandReply{Text: "用法:<code>/cleanup_rule del 规则ID</code>"}
}
next := make([]accountCleanupRule, 0, len(rules))
removed := false
for _, r := range rules {
if r.ID == args[1] {
removed = true
continue
}
next = append(next, r)
}
if !removed {
return telegramCommandReply{Text: "未找到该规则。"}
}
if err := s.saveCleanupRules(ctx, next); err != nil {
return telegramCommandReply{Text: "保存失败:" + err.Error()}
}
return telegramCommandReply{Text: "已删除规则。\n\n" + formatCleanupRules(next)}
case "enable", "on", "disable", "off":
if len(args) < 2 {
return telegramCommandReply{Text: "用法:<code>/cleanup_rule enable|disable 规则ID</code>"}
}
enable := action == "enable" || action == "on"
changed := false
for i := range rules {
if rules[i].ID == args[1] {
rules[i].Enabled = enable
changed = true
}
}
if !changed {
return telegramCommandReply{Text: "未找到该规则。"}
}
if err := s.saveCleanupRules(ctx, rules); err != nil {
return telegramCommandReply{Text: "保存失败:" + err.Error()}
}
return telegramCommandReply{Text: "已更新规则状态。\n\n" + formatCleanupRules(rules)}
case "add":
rule, err := parseCleanupRuleCommand(args[1:])
if err != nil {
return telegramCommandReply{Text: err.Error() + "\n\n" + cleanupRuleHelp()}
}
for _, r := range rules {
if r.ID == rule.ID {
return telegramCommandReply{Text: "规则 ID 已存在,请换一个 ID。"}
}
}
rules = normalizeCleanupRules(append(rules, rule))
if err := s.saveCleanupRules(ctx, rules); err != nil {
return telegramCommandReply{Text: "保存失败:" + err.Error()}
}
return telegramCommandReply{Text: "已新增规则。\n\n" + formatCleanupRules(rules)}
default:
return telegramCommandReply{Text: cleanupRuleHelp()}
}
}
func (s *TelegramBotService) replyDevicePolicyToggle(ctx context.Context, which string) telegramCommandReply {
cfg := loadBotConfig(ctx, s.repo)
switch which {
@@ -607,6 +827,139 @@ func (s *TelegramBotService) replyDevicePolicyToggle(ctx context.Context, which
return s.replyDevicePolicy(ctx)
}
func (s *TelegramBotService) cmdUserBan(ctx context.Context, args []string, unban bool) telegramCommandReply {
if len(args) == 0 {
if unban {
return telegramCommandReply{Text: "用法:<code>/unban 用户名</code>"}
}
return telegramCommandReply{Text: "用法:<code>/ban 用户名</code>"}
}
user, _ := s.repo.User.FindByUsername(ctx, args[0])
if user == nil {
user, _ = s.repo.User.FindByID(ctx, args[0])
}
if user == nil {
return telegramCommandReply{Text: "未找到用户。"}
}
return s.replyUserBan(ctx, user.ID, unban)
}
func (s *TelegramBotService) currentCleanupRules(ctx context.Context) []accountCleanupRule {
cfg := loadBotConfig(ctx, s.repo)
return cfg.AccountCleanupRules
}
func (s *TelegramBotService) saveCleanupRules(ctx context.Context, rules []accountCleanupRule) error {
raw, err := json.Marshal(normalizeCleanupRules(rules))
if err != nil {
return err
}
return s.repo.Setting.Set(ctx, SettingAccountCleanupRules, string(raw))
}
func parseCommandBool(value string) (bool, bool) {
switch strings.ToLower(strings.TrimSpace(value)) {
case "on", "true", "1", "yes", "enable", "enabled", "开启", "开":
return true, true
case "off", "false", "0", "no", "disable", "disabled", "关闭", "关":
return false, true
default:
return false, false
}
}
func parseCleanupRuleCommand(args []string) (accountCleanupRule, error) {
if len(args) < 2 {
return accountCleanupRule{}, fmt.Errorf("新增规则参数不足")
}
rule := accountCleanupRule{
Type: strings.ToLower(strings.TrimSpace(args[0])),
ID: strings.TrimSpace(args[1]),
Name: strings.TrimSpace(args[1]),
Enabled: true,
WindowDaysMin: 3,
WindowDaysMax: 5,
MinHours: 6,
MinCount: 1,
}
if len(args) > 2 {
rule.Name = strings.TrimSpace(args[2])
}
switch rule.Type {
case "watch_hours":
if len(args) >= 6 {
rule.WindowDaysMin, _ = strconv.Atoi(args[3])
rule.WindowDaysMax, _ = strconv.Atoi(args[4])
rule.MinHours, _ = strconv.ParseFloat(args[5], 64)
}
case "recent_login":
if len(args) >= 4 {
rule.WindowDaysMax, _ = strconv.Atoi(args[3])
}
case "signin_streak", "account_age_grace":
if len(args) >= 4 {
rule.MinCount, _ = strconv.Atoi(args[3])
}
default:
return accountCleanupRule{}, fmt.Errorf("不支持的规则类型:%s", rule.Type)
}
normalized := normalizeCleanupRules([]accountCleanupRule{rule})
if len(normalized) == 0 {
return accountCleanupRule{}, fmt.Errorf("规则无效")
}
return normalized[0], nil
}
func formatCleanupRules(rules []accountCleanupRule) string {
if len(rules) == 0 {
return "<b>保号规则</b>\n\n暂无规则。"
}
var sb strings.Builder
sb.WriteString("<b>保号规则</b>\n")
for i, r := range rules {
state := map[bool]string{true: "启用", false: "停用"}[r.Enabled]
sb.WriteString(fmt.Sprintf("\n%d. <code>%s</code> · %s · %s · %s", i+1, r.ID, r.Name, cleanupRuleTypeLabel(r.Type), state))
switch r.Type {
case "watch_hours":
sb.WriteString(fmt.Sprintf(" · %d~%d 天 %.1f 小时", r.WindowDaysMin, r.WindowDaysMax, r.MinHours))
case "recent_login":
sb.WriteString(fmt.Sprintf(" · %d 天内登录", r.WindowDaysMax))
case "signin_streak":
sb.WriteString(fmt.Sprintf(" · 连续签到 %d 天", r.MinCount))
case "account_age_grace":
sb.WriteString(fmt.Sprintf(" · 新号宽限 %d 天", r.MinCount))
}
}
return sb.String()
}
func cleanupRuleTypeLabel(t string) string {
switch t {
case "watch_hours":
return "观看时长"
case "recent_login":
return "最近登录"
case "signin_streak":
return "连续签到"
case "account_age_grace":
return "新号宽限"
default:
return t
}
}
func cleanupRuleHelp() string {
return "<b>删号/保号规则命令</b>\n\n" +
"<code>/cleanup_rule list</code> — 查看规则\n" +
"<code>/cleanup_rule add watch_hours watch_3_5d_6h 观看3到5天满6小时 3 5 6</code>\n" +
"<code>/cleanup_rule add recent_login login_7d 七天内登录 7</code>\n" +
"<code>/cleanup_rule add signin_streak sign_3 连续签到3天 3</code>\n" +
"<code>/cleanup_rule add account_age_grace new_7d 新号宽限7天 7</code>\n" +
"<code>/cleanup_rule enable 规则ID</code> / <code>disable 规则ID</code>\n" +
"<code>/cleanup_rule del 规则ID</code>\n\n" +
"保号模式:<code>/cleanup_mode any|all|count 2</code>"
}
func onOff(b bool) string {
return map[bool]string{true: "已开启", false: "已关闭"}[b]
}
+11
View File
@@ -84,6 +84,9 @@ func (s *TokenService) IssuePair(ctx context.Context, userID, role, tier string)
if err := s.repo.RefreshToken.Create(ctx, rt); err != nil {
return nil, err
}
if err := s.repo.RefreshToken.RevokeOldestActiveByUserID(ctx, userID, s.maxActiveRefreshTokens(ctx)); err != nil {
s.log.Warn("failed to enforce refresh token session limit", zap.String("user_id", userID), zap.Error(err))
}
return &TokenPair{
AccessToken: accessToken,
@@ -93,6 +96,14 @@ func (s *TokenService) IssuePair(ctx context.Context, userID, role, tier string)
}, nil
}
func (s *TokenService) maxActiveRefreshTokens(ctx context.Context) int {
cfg := loadBotConfig(ctx, s.repo)
if cfg.MaxLoggedClients < 1 {
return defaultBotConfig().MaxLoggedClients
}
return cfg.MaxLoggedClients
}
// issueAccessToken 签发 JWT Access Token(HS256,60分钟有效期)。
func (s *TokenService) issueAccessToken(userID, role, tier string) (string, error) {
claims := Claims{