From cec89cad3bb2ffab0bc05f0d4a0f66f535367559 Mon Sep 17 00:00:00 2001
From: ShukeBta <272197458+ShukeBta@users.noreply.github.com>
Date: Fri, 26 Jun 2026 18:32:50 +0800
Subject: [PATCH] refactor telegram mgo commands by topic
---
internal/service/telegram_mgo_admin.go | 159 +++++++
internal/service/telegram_mgo_audit.go | 133 ++++++
internal/service/telegram_mgo_backup.go | 49 ++
internal/service/telegram_mgo_bulk.go | 94 ++++
internal/service/telegram_mgo_compat.go | 481 --------------------
internal/service/telegram_mgo_membership.go | 92 ++++
6 files changed, 527 insertions(+), 481 deletions(-)
create mode 100644 internal/service/telegram_mgo_admin.go
create mode 100644 internal/service/telegram_mgo_audit.go
create mode 100644 internal/service/telegram_mgo_backup.go
create mode 100644 internal/service/telegram_mgo_bulk.go
create mode 100644 internal/service/telegram_mgo_membership.go
diff --git a/internal/service/telegram_mgo_admin.go b/internal/service/telegram_mgo_admin.go
new file mode 100644
index 0000000..728d50f
--- /dev/null
+++ b/internal/service/telegram_mgo_admin.go
@@ -0,0 +1,159 @@
+package service
+
+import (
+ "context"
+ "encoding/json"
+ "fmt"
+ "sort"
+ "strconv"
+ "strings"
+
+ "github.com/ShukeBta/MediaStationGo/internal/model"
+)
+
+func (s *TelegramBotService) cmdMgoAdminRole(ctx context.Context, args []string) telegramCommandReply {
+ if len(args) < 2 {
+ return telegramCommandReply{Text: "用法:/embyadmin 用户名 on|off"}
+ }
+ user := s.findMgoBotUser(ctx, args[0])
+ if user == nil {
+ return telegramCommandReply{Text: "未找到用户。"}
+ }
+ enable := parseOnOff(args[1])
+ if enable == nil {
+ return telegramCommandReply{Text: "第二个参数请使用 on/off。"}
+ }
+ if !*enable {
+ if first, _ := s.repo.User.FirstAdmin(ctx); first != nil && first.ID == user.ID {
+ return telegramCommandReply{Text: "默认管理员不可降级。"}
+ }
+ }
+ role := "user"
+ if *enable {
+ role = "admin"
+ }
+ if err := s.repo.User.UpdateFields(ctx, user.ID, map[string]any{"role": role}); err != nil {
+ return telegramCommandReply{Text: "更新失败:" + err.Error()}
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("已将 %s 角色设置为 %s。", user.Username, role)}
+}
+
+func (s *TelegramBotService) cmdMgoMediaAccessAll(ctx context.Context, allow bool) telegramCommandReply {
+ users, err := s.repo.User.List(ctx)
+ if err != nil {
+ return telegramCommandReply{Text: "读取用户失败:" + err.Error()}
+ }
+ updated := 0
+ for _, user := range users {
+ if user.Role == "admin" {
+ continue
+ }
+ perm, err := s.repo.Permission.FindByUserID(ctx, user.ID)
+ if err != nil {
+ continue
+ }
+ if perm == nil {
+ perm = DefaultPermissions(user.ID)
+ perm.CanPlayMedia = allow
+ if err := s.repo.Permission.Create(ctx, perm); err != nil {
+ continue
+ }
+ }
+ if err := s.repo.DB.WithContext(ctx).Model(&model.UserPermission{}).
+ Where("user_id = ?", user.ID).
+ Update("can_play_media", allow).Error; err == nil {
+ updated++
+ }
+ }
+ state := "关闭"
+ if allow {
+ state = "开启"
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("已为普通用户%s媒体播放权限:%d 个。", state, updated)}
+}
+
+func (s *TelegramBotService) cmdMgoBotAdmin(ctx context.Context, channel *model.NotifyChannel, args []string, add bool) telegramCommandReply {
+ if channel == nil {
+ return telegramCommandReply{Text: "Telegram 渠道不存在。"}
+ }
+ if len(args) == 0 {
+ return telegramCommandReply{Text: "用法:/proadmin TelegramID 或 /revadmin TelegramID"}
+ }
+ tgID := strings.TrimPrefix(strings.TrimSpace(args[0]), "tg:")
+ if _, err := strconv.ParseInt(tgID, 10, 64); err != nil {
+ return telegramCommandReply{Text: "TelegramID 必须是数字。"}
+ }
+ cfg := s.telegramChannelConfig(channel)
+ ids := telegramConfiguredUserIDs(cfg["admin_user_ids"])
+ seen := make(map[string]bool, len(ids)+1)
+ var next []string
+ for _, id := range ids {
+ if id == tgID {
+ seen[id] = true
+ if add {
+ next = append(next, id)
+ }
+ continue
+ }
+ if id != "" {
+ next = append(next, id)
+ }
+ }
+ if add && !seen[tgID] {
+ next = append(next, tgID)
+ }
+ cfg["admin_user_ids"] = strings.Join(next, ",")
+ raw, _ := json.Marshal(cfg)
+ updated := *channel
+ updated.Config = string(raw)
+ if s.crypto != nil {
+ updated.Config = s.crypto.Encrypt(updated.Config)
+ }
+ if err := s.repo.NotifyChannel.Update(ctx, &updated); err != nil {
+ return telegramCommandReply{Text: "更新管理员列表失败:" + err.Error()}
+ }
+ if add {
+ return telegramCommandReply{Text: "已添加 Bot 管理员:" + tgID + ""}
+ }
+ return telegramCommandReply{Text: "已移除 Bot 管理员:" + tgID + ""}
+}
+
+func (s *TelegramBotService) cmdMgoProtectedUser(ctx context.Context, args []string, protect bool) telegramCommandReply {
+ if len(args) == 0 || strings.EqualFold(args[0], "list") {
+ return s.cmdMgoProtectedUserList(ctx)
+ }
+ user := s.findMgoBotUser(ctx, args[0])
+ if user == nil {
+ return telegramCommandReply{Text: "未找到用户。"}
+ }
+ ids := ProtectedUserIDSet(ctx, s.repo)
+ if protect {
+ ids[user.ID] = struct{}{}
+ if err := SaveProtectedUserIDSet(ctx, s.repo, ids); err != nil {
+ return telegramCommandReply{Text: "保存保护名单失败:" + err.Error()}
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("已加入保护名单:%s。\n该用户不会被 Bot 自动清理、批量禁用或删除。", user.Username)}
+ }
+ delete(ids, user.ID)
+ if err := SaveProtectedUserIDSet(ctx, s.repo, ids); err != nil {
+ return telegramCommandReply{Text: "保存保护名单失败:" + err.Error()}
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("已移出保护名单:%s。", user.Username)}
+}
+
+func (s *TelegramBotService) cmdMgoProtectedUserList(ctx context.Context) telegramCommandReply {
+ ids := ProtectedUserIDSet(ctx, s.repo)
+ if len(ids) == 0 {
+ return telegramCommandReply{Text: "保护名单为空。管理员和默认管理员始终自动保护。"}
+ }
+ names := make([]string, 0, len(ids))
+ for id := range ids {
+ if user, _ := s.repo.User.FindByID(ctx, id); user != nil {
+ names = append(names, user.Username)
+ } else {
+ names = append(names, id+"(用户不存在)")
+ }
+ }
+ sort.Strings(names)
+ return telegramCommandReply{Text: fmt.Sprintf("保护名单:%d 个。\n%s", len(names), telegramInlineCodeList(names))}
+}
diff --git a/internal/service/telegram_mgo_audit.go b/internal/service/telegram_mgo_audit.go
new file mode 100644
index 0000000..b696852
--- /dev/null
+++ b/internal/service/telegram_mgo_audit.go
@@ -0,0 +1,133 @@
+package service
+
+import (
+ "context"
+ "fmt"
+ "strings"
+ "time"
+
+ "github.com/ShukeBta/MediaStationGo/internal/model"
+)
+
+func (s *TelegramBotService) cmdMgoSyncUnbound(ctx context.Context, args []string) telegramCommandReply {
+ var users []model.User
+ if err := s.repo.DB.WithContext(ctx).
+ Where("role <> ?", "admin").
+ Where("NOT EXISTS (SELECT 1 FROM telegram_bindings WHERE telegram_bindings.user_id = users.id AND telegram_bindings.deleted_at IS NULL)").
+ Order("created_at asc").Find(&users).Error; err != nil {
+ return telegramCommandReply{Text: "查询失败:" + err.Error()}
+ }
+ if len(args) >= 2 && strings.EqualFold(args[0], "delete") && strings.EqualFold(args[1], "confirm") {
+ deleted := 0
+ for _, user := range users {
+ if UserIsProtectedAccount(ctx, s.repo, &user) {
+ continue
+ }
+ _ = s.repo.UserDevice.DeleteByUser(ctx, user.ID)
+ if err := s.repo.User.Delete(ctx, user.ID); err == nil {
+ deleted++
+ }
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("已删除未绑定 Bot 的普通用户:%d 个。", deleted)}
+ }
+ if len(users) == 0 {
+ return telegramCommandReply{Text: "没有未绑定 Bot 的普通用户。"}
+ }
+ names := make([]string, 0, minInt(len(users), 20))
+ for i, user := range users {
+ if i >= 20 {
+ break
+ }
+ names = append(names, user.Username)
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("未绑定 Bot 的普通用户:%d 个。\n%s\n\n如需删除:/syncunbound delete confirm", len(users), telegramInlineCodeList(names))}
+}
+
+func (s *TelegramBotService) cmdMgoCheckExpired(ctx context.Context, args []string) telegramCommandReply {
+ now := time.Now()
+ var users []model.User
+ if err := s.repo.DB.WithContext(ctx).Where("expired_at IS NOT NULL AND expired_at < ?", now).Order("expired_at asc").Find(&users).Error; err != nil {
+ return telegramCommandReply{Text: "查询失败:" + err.Error()}
+ }
+ if len(args) >= 2 && strings.EqualFold(args[0], "disable") && strings.EqualFold(args[1], "confirm") {
+ disabled := 0
+ for _, user := range users {
+ if UserIsProtectedAccount(ctx, s.repo, &user) {
+ continue
+ }
+ if err := s.repo.User.UpdateFields(ctx, user.ID, map[string]any{"is_active": false}); err == nil {
+ disabled++
+ }
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("已禁用过期普通用户:%d 个。", disabled)}
+ }
+ if len(users) == 0 {
+ return telegramCommandReply{Text: "没有过期用户。"}
+ }
+ lines := make([]string, 0, minInt(len(users), 20))
+ for i, user := range users {
+ if i >= 20 {
+ break
+ }
+ lines = append(lines, fmt.Sprintf("%s(%s)", user.Username, formatExpiry(user.ExpiredAt)))
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("过期用户:%d 个。\n%s\n\n如需禁用:/check_ex disable confirm", len(users), telegramInlineCodeList(lines))}
+}
+
+func (s *TelegramBotService) cmdMgoScanNames(ctx context.Context) telegramCommandReply {
+ var rows []struct {
+ Username string
+ Count int64
+ }
+ if err := s.repo.DB.WithContext(ctx).Table("users").
+ Select("LOWER(username) AS username, COUNT(*) AS count").
+ Group("LOWER(username)").Having("COUNT(*) > 1").Scan(&rows).Error; err != nil {
+ return telegramCommandReply{Text: "扫描失败:" + err.Error()}
+ }
+ if len(rows) == 0 {
+ return telegramCommandReply{Text: "未发现同名用户记录。"}
+ }
+ var out []string
+ for _, row := range rows {
+ out = append(out, fmt.Sprintf("%s x%d", row.Username, row.Count))
+ }
+ return telegramCommandReply{Text: "同名用户记录\n" + telegramInlineCodeList(out)}
+}
+
+func (s *TelegramBotService) cmdMgoRanks(ctx context.Context, window time.Duration, byDuration bool) telegramCommandReply {
+ since := time.Now().Add(-window)
+ title := "播放次数排行"
+ selectExpr := "COUNT(*) AS score"
+ if byDuration {
+ title = "观影时长排行"
+ selectExpr = "COALESCE(SUM(position_ms), 0) AS score"
+ }
+ q := s.repo.DB.WithContext(ctx).Table("playback_histories").
+ Select("users.username, " + selectExpr).
+ Joins("JOIN users ON users.id = playback_histories.user_id").
+ Group("users.username").
+ Order("score DESC").
+ Limit(10)
+ if window > 0 {
+ q = q.Where("playback_histories.watched_at >= ?", since)
+ }
+ var rows []struct {
+ Username string
+ Score int64
+ }
+ if err := q.Scan(&rows).Error; err != nil {
+ return telegramCommandReply{Text: "排行查询失败:" + err.Error()}
+ }
+ if len(rows) == 0 {
+ return telegramCommandReply{Text: "暂无排行数据。"}
+ }
+ var out []string
+ for i, row := range rows {
+ score := fmt.Sprintf("%d 次", row.Score)
+ if byDuration {
+ score = humanDurationFromMillis(row.Score)
+ }
+ out = append(out, fmt.Sprintf("%d. %s — %s", i+1, row.Username, score))
+ }
+ return telegramCommandReply{Text: "" + title + "\n\n" + strings.Join(out, "\n") + ""}
+}
diff --git a/internal/service/telegram_mgo_backup.go b/internal/service/telegram_mgo_backup.go
new file mode 100644
index 0000000..7ba58c9
--- /dev/null
+++ b/internal/service/telegram_mgo_backup.go
@@ -0,0 +1,49 @@
+package service
+
+import (
+ "context"
+ "fmt"
+ "strings"
+)
+
+func (s *TelegramBotService) cmdMgoBackupDB(ctx context.Context) telegramCommandReply {
+ if s.backup == nil {
+ return telegramCommandReply{Text: "备份服务暂不可用。"}
+ }
+ info, err := s.backup.Create(ctx)
+ if err != nil {
+ return telegramCommandReply{Text: "数据库备份失败:" + err.Error()}
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("数据库备份完成:%s\n大小:%d bytes", info.Filename, info.Size)}
+}
+
+func (s *TelegramBotService) cmdMgoRestoreDB(ctx context.Context, args []string) telegramCommandReply {
+ if s.backup == nil {
+ return telegramCommandReply{Text: "备份服务暂不可用。"}
+ }
+ if len(args) == 0 || strings.EqualFold(args[0], "list") {
+ items, err := s.backup.List()
+ if err != nil {
+ return telegramCommandReply{Text: "读取备份列表失败:" + err.Error()}
+ }
+ if len(items) == 0 {
+ return telegramCommandReply{Text: "暂无数据库备份。可先使用 /backup_db 创建。"}
+ }
+ lines := make([]string, 0, minInt(len(items), 10))
+ for i, item := range items {
+ if i >= 10 {
+ break
+ }
+ lines = append(lines, fmt.Sprintf("%s(%d bytes)", item.Filename, item.Size))
+ }
+ return telegramCommandReply{Text: "可恢复备份:\n" + telegramInlineCodeList(lines) + "\n\n恢复需要确认:/restore_from_db 文件名 confirm"}
+ }
+ if len(args) < 2 || !strings.EqualFold(args[len(args)-1], "confirm") {
+ return telegramCommandReply{Text: "恢复数据库会覆盖当前数据,需要确认:/restore_from_db 文件名 confirm"}
+ }
+ filename := strings.TrimSpace(args[0])
+ if err := s.backup.Restore(ctx, filename); err != nil {
+ return telegramCommandReply{Text: "恢复失败:" + err.Error()}
+ }
+ return telegramCommandReply{Text: "数据库已从备份恢复,请重启 MediaStationGo 后生效。"}
+}
diff --git a/internal/service/telegram_mgo_bulk.go b/internal/service/telegram_mgo_bulk.go
new file mode 100644
index 0000000..e06f554
--- /dev/null
+++ b/internal/service/telegram_mgo_bulk.go
@@ -0,0 +1,94 @@
+package service
+
+import (
+ "context"
+ "fmt"
+ "strconv"
+ "strings"
+
+ "github.com/ShukeBta/MediaStationGo/internal/model"
+)
+
+func (s *TelegramBotService) cmdMgoRenewAll(ctx context.Context, args []string) telegramCommandReply {
+ if len(args) < 2 || !strings.EqualFold(args[len(args)-1], "confirm") {
+ return telegramCommandReply{Text: "批量续期需要确认:/renewall 天数 confirm"}
+ }
+ days, err := strconv.Atoi(args[0])
+ if err != nil || days < 0 {
+ return telegramCommandReply{Text: "天数必须是非负整数,0 表示永久。"}
+ }
+ users, err := s.repo.User.List(ctx)
+ if err != nil {
+ return telegramCommandReply{Text: "读取用户失败:" + err.Error()}
+ }
+ var count int
+ for _, user := range users {
+ if user.Role == "admin" {
+ continue
+ }
+ if err := s.applyRenewal(ctx, user.ID, days); err == nil {
+ count++
+ }
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("批量续期完成:%d 个普通用户。", count)}
+}
+
+func (s *TelegramBotService) cmdMgoBanAll(ctx context.Context, active bool, args []string) telegramCommandReply {
+ if len(args) == 0 || !strings.EqualFold(args[len(args)-1], "confirm") {
+ action := "banall"
+ if active {
+ action = "unbanall"
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("批量操作需要确认:/%s confirm", action)}
+ }
+ users, err := s.repo.User.List(ctx)
+ if err != nil {
+ return telegramCommandReply{Text: "读取用户失败:" + err.Error()}
+ }
+ var count int
+ for _, user := range users {
+ if !active && UserIsProtectedAccount(ctx, s.repo, &user) {
+ continue
+ }
+ if active && user.Role == "admin" {
+ continue
+ }
+ updates := map[string]any{"is_active": active}
+ if active {
+ updates["share_warnings"] = 0
+ updates["last_share_warn_at"] = nil
+ }
+ if err := s.repo.User.UpdateFields(ctx, user.ID, updates); err == nil {
+ _ = s.repo.UserDevice.SetKickedByUser(ctx, user.ID, !active)
+ count++
+ }
+ }
+ if active {
+ return telegramCommandReply{Text: fmt.Sprintf("已解禁普通用户:%d 个。", count)}
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("已禁用普通用户:%d 个。", count)}
+}
+
+func (s *TelegramBotService) cmdMgoCallAll(ctx context.Context, channel *model.NotifyChannel, args []string) telegramCommandReply {
+ message := strings.TrimSpace(strings.Join(args, " "))
+ if message == "" {
+ return telegramCommandReply{Text: "用法:/callall 消息内容"}
+ }
+ if strings.TrimSpace(s.telegramChannelConfig(channel)["bot_token"]) == "" {
+ return telegramCommandReply{Text: "当前 Telegram 渠道未配置 bot_token,无法群发。"}
+ }
+ var bindings []model.TelegramBinding
+ if err := s.repo.DB.WithContext(ctx).Find(&bindings).Error; err != nil {
+ return telegramCommandReply{Text: "读取绑定失败:" + err.Error()}
+ }
+ sent := 0
+ for _, binding := range bindings {
+ if binding.ChatID == 0 {
+ continue
+ }
+ if err := s.reply(ctx, channel, int(binding.ChatID), telegramCommandReply{Text: message}); err == nil {
+ sent++
+ }
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("群发完成:成功发送 %d 个绑定用户。", sent)}
+}
diff --git a/internal/service/telegram_mgo_compat.go b/internal/service/telegram_mgo_compat.go
index bd5a928..46f2188 100644
--- a/internal/service/telegram_mgo_compat.go
+++ b/internal/service/telegram_mgo_compat.go
@@ -2,9 +2,7 @@ package service
import (
"context"
- "encoding/json"
"fmt"
- "sort"
"strconv"
"strings"
"time"
@@ -12,485 +10,6 @@ import (
"github.com/ShukeBta/MediaStationGo/internal/model"
)
-func (s *TelegramBotService) cmdMgoRenewAll(ctx context.Context, args []string) telegramCommandReply {
- if len(args) < 2 || !strings.EqualFold(args[len(args)-1], "confirm") {
- return telegramCommandReply{Text: "批量续期需要确认:/renewall 天数 confirm"}
- }
- days, err := strconv.Atoi(args[0])
- if err != nil || days < 0 {
- return telegramCommandReply{Text: "天数必须是非负整数,0 表示永久。"}
- }
- users, err := s.repo.User.List(ctx)
- if err != nil {
- return telegramCommandReply{Text: "读取用户失败:" + err.Error()}
- }
- var count int
- for _, user := range users {
- if user.Role == "admin" {
- continue
- }
- if err := s.applyRenewal(ctx, user.ID, days); err == nil {
- count++
- }
- }
- return telegramCommandReply{Text: fmt.Sprintf("批量续期完成:%d 个普通用户。", count)}
-}
-
-func (s *TelegramBotService) cmdMgoBanAll(ctx context.Context, active bool, args []string) telegramCommandReply {
- if len(args) == 0 || !strings.EqualFold(args[len(args)-1], "confirm") {
- action := "banall"
- if active {
- action = "unbanall"
- }
- return telegramCommandReply{Text: fmt.Sprintf("批量操作需要确认:/%s confirm", action)}
- }
- users, err := s.repo.User.List(ctx)
- if err != nil {
- return telegramCommandReply{Text: "读取用户失败:" + err.Error()}
- }
- var count int
- for _, user := range users {
- if !active && UserIsProtectedAccount(ctx, s.repo, &user) {
- continue
- }
- if active && user.Role == "admin" {
- continue
- }
- updates := map[string]any{"is_active": active}
- if active {
- updates["share_warnings"] = 0
- updates["last_share_warn_at"] = nil
- }
- if err := s.repo.User.UpdateFields(ctx, user.ID, updates); err == nil {
- _ = s.repo.UserDevice.SetKickedByUser(ctx, user.ID, !active)
- count++
- }
- }
- if active {
- return telegramCommandReply{Text: fmt.Sprintf("已解禁普通用户:%d 个。", count)}
- }
- return telegramCommandReply{Text: fmt.Sprintf("已禁用普通用户:%d 个。", count)}
-}
-
-func (s *TelegramBotService) cmdMgoCallAll(ctx context.Context, channel *model.NotifyChannel, args []string) telegramCommandReply {
- message := strings.TrimSpace(strings.Join(args, " "))
- if message == "" {
- return telegramCommandReply{Text: "用法:/callall 消息内容"}
- }
- if strings.TrimSpace(s.telegramChannelConfig(channel)["bot_token"]) == "" {
- return telegramCommandReply{Text: "当前 Telegram 渠道未配置 bot_token,无法群发。"}
- }
- var bindings []model.TelegramBinding
- if err := s.repo.DB.WithContext(ctx).Find(&bindings).Error; err != nil {
- return telegramCommandReply{Text: "读取绑定失败:" + err.Error()}
- }
- sent := 0
- for _, binding := range bindings {
- if binding.ChatID == 0 {
- continue
- }
- if err := s.reply(ctx, channel, int(binding.ChatID), telegramCommandReply{Text: message}); err == nil {
- sent++
- }
- }
- return telegramCommandReply{Text: fmt.Sprintf("群发完成:成功发送 %d 个绑定用户。", sent)}
-}
-
-func (s *TelegramBotService) cmdMgoSyncUnbound(ctx context.Context, args []string) telegramCommandReply {
- var users []model.User
- if err := s.repo.DB.WithContext(ctx).
- Where("role <> ?", "admin").
- Where("NOT EXISTS (SELECT 1 FROM telegram_bindings WHERE telegram_bindings.user_id = users.id AND telegram_bindings.deleted_at IS NULL)").
- Order("created_at asc").Find(&users).Error; err != nil {
- return telegramCommandReply{Text: "查询失败:" + err.Error()}
- }
- if len(args) >= 2 && strings.EqualFold(args[0], "delete") && strings.EqualFold(args[1], "confirm") {
- deleted := 0
- for _, user := range users {
- if UserIsProtectedAccount(ctx, s.repo, &user) {
- continue
- }
- _ = s.repo.UserDevice.DeleteByUser(ctx, user.ID)
- if err := s.repo.User.Delete(ctx, user.ID); err == nil {
- deleted++
- }
- }
- return telegramCommandReply{Text: fmt.Sprintf("已删除未绑定 Bot 的普通用户:%d 个。", deleted)}
- }
- if len(users) == 0 {
- return telegramCommandReply{Text: "没有未绑定 Bot 的普通用户。"}
- }
- names := make([]string, 0, minInt(len(users), 20))
- for i, user := range users {
- if i >= 20 {
- break
- }
- names = append(names, user.Username)
- }
- return telegramCommandReply{Text: fmt.Sprintf("未绑定 Bot 的普通用户:%d 个。\n%s\n\n如需删除:/syncunbound delete confirm", len(users), telegramInlineCodeList(names))}
-}
-
-func (s *TelegramBotService) cmdMgoCheckExpired(ctx context.Context, args []string) telegramCommandReply {
- now := time.Now()
- var users []model.User
- if err := s.repo.DB.WithContext(ctx).Where("expired_at IS NOT NULL AND expired_at < ?", now).Order("expired_at asc").Find(&users).Error; err != nil {
- return telegramCommandReply{Text: "查询失败:" + err.Error()}
- }
- if len(args) >= 2 && strings.EqualFold(args[0], "disable") && strings.EqualFold(args[1], "confirm") {
- disabled := 0
- for _, user := range users {
- if UserIsProtectedAccount(ctx, s.repo, &user) {
- continue
- }
- if err := s.repo.User.UpdateFields(ctx, user.ID, map[string]any{"is_active": false}); err == nil {
- disabled++
- }
- }
- return telegramCommandReply{Text: fmt.Sprintf("已禁用过期普通用户:%d 个。", disabled)}
- }
- if len(users) == 0 {
- return telegramCommandReply{Text: "没有过期用户。"}
- }
- lines := make([]string, 0, minInt(len(users), 20))
- for i, user := range users {
- if i >= 20 {
- break
- }
- lines = append(lines, fmt.Sprintf("%s(%s)", user.Username, formatExpiry(user.ExpiredAt)))
- }
- return telegramCommandReply{Text: fmt.Sprintf("过期用户:%d 个。\n%s\n\n如需禁用:/check_ex disable confirm", len(users), telegramInlineCodeList(lines))}
-}
-
-func (s *TelegramBotService) cmdMgoScanNames(ctx context.Context) telegramCommandReply {
- var rows []struct {
- Username string
- Count int64
- }
- if err := s.repo.DB.WithContext(ctx).Table("users").
- Select("LOWER(username) AS username, COUNT(*) AS count").
- Group("LOWER(username)").Having("COUNT(*) > 1").Scan(&rows).Error; err != nil {
- return telegramCommandReply{Text: "扫描失败:" + err.Error()}
- }
- if len(rows) == 0 {
- return telegramCommandReply{Text: "未发现同名用户记录。"}
- }
- var out []string
- for _, row := range rows {
- out = append(out, fmt.Sprintf("%s x%d", row.Username, row.Count))
- }
- return telegramCommandReply{Text: "同名用户记录\n" + telegramInlineCodeList(out)}
-}
-
-func (s *TelegramBotService) cmdMgoRanks(ctx context.Context, window time.Duration, byDuration bool) telegramCommandReply {
- since := time.Now().Add(-window)
- title := "播放次数排行"
- selectExpr := "COUNT(*) AS score"
- if byDuration {
- title = "观影时长排行"
- selectExpr = "COALESCE(SUM(position_ms), 0) AS score"
- }
- q := s.repo.DB.WithContext(ctx).Table("playback_histories").
- Select("users.username, " + selectExpr).
- Joins("JOIN users ON users.id = playback_histories.user_id").
- Group("users.username").
- Order("score DESC").
- Limit(10)
- if window > 0 {
- q = q.Where("playback_histories.watched_at >= ?", since)
- }
- var rows []struct {
- Username string
- Score int64
- }
- if err := q.Scan(&rows).Error; err != nil {
- return telegramCommandReply{Text: "排行查询失败:" + err.Error()}
- }
- if len(rows) == 0 {
- return telegramCommandReply{Text: "暂无排行数据。"}
- }
- var out []string
- for i, row := range rows {
- score := fmt.Sprintf("%d 次", row.Score)
- if byDuration {
- score = humanDurationFromMillis(row.Score)
- }
- out = append(out, fmt.Sprintf("%d. %s — %s", i+1, row.Username, score))
- }
- return telegramCommandReply{Text: "" + title + "\n\n" + strings.Join(out, "\n") + ""}
-}
-
-func (s *TelegramBotService) cmdMgoAdminRole(ctx context.Context, args []string) telegramCommandReply {
- if len(args) < 2 {
- return telegramCommandReply{Text: "用法:/embyadmin 用户名 on|off"}
- }
- user := s.findMgoBotUser(ctx, args[0])
- if user == nil {
- return telegramCommandReply{Text: "未找到用户。"}
- }
- enable := parseOnOff(args[1])
- if enable == nil {
- return telegramCommandReply{Text: "第二个参数请使用 on/off。"}
- }
- if !*enable {
- if first, _ := s.repo.User.FirstAdmin(ctx); first != nil && first.ID == user.ID {
- return telegramCommandReply{Text: "默认管理员不可降级。"}
- }
- }
- role := "user"
- if *enable {
- role = "admin"
- }
- if err := s.repo.User.UpdateFields(ctx, user.ID, map[string]any{"role": role}); err != nil {
- return telegramCommandReply{Text: "更新失败:" + err.Error()}
- }
- return telegramCommandReply{Text: fmt.Sprintf("已将 %s 角色设置为 %s。", user.Username, role)}
-}
-
-func (s *TelegramBotService) cmdMgoMediaAccessAll(ctx context.Context, allow bool) telegramCommandReply {
- users, err := s.repo.User.List(ctx)
- if err != nil {
- return telegramCommandReply{Text: "读取用户失败:" + err.Error()}
- }
- updated := 0
- for _, user := range users {
- if user.Role == "admin" {
- continue
- }
- perm, err := s.repo.Permission.FindByUserID(ctx, user.ID)
- if err != nil {
- continue
- }
- if perm == nil {
- perm = DefaultPermissions(user.ID)
- perm.CanPlayMedia = allow
- if err := s.repo.Permission.Create(ctx, perm); err != nil {
- continue
- }
- }
- if err := s.repo.DB.WithContext(ctx).Model(&model.UserPermission{}).
- Where("user_id = ?", user.ID).
- Update("can_play_media", allow).Error; err == nil {
- updated++
- }
- }
- state := "关闭"
- if allow {
- state = "开启"
- }
- return telegramCommandReply{Text: fmt.Sprintf("已为普通用户%s媒体播放权限:%d 个。", state, updated)}
-}
-
-func (s *TelegramBotService) cmdMgoBotAdmin(ctx context.Context, channel *model.NotifyChannel, args []string, add bool) telegramCommandReply {
- if channel == nil {
- return telegramCommandReply{Text: "Telegram 渠道不存在。"}
- }
- if len(args) == 0 {
- return telegramCommandReply{Text: "用法:/proadmin TelegramID 或 /revadmin TelegramID"}
- }
- tgID := strings.TrimPrefix(strings.TrimSpace(args[0]), "tg:")
- if _, err := strconv.ParseInt(tgID, 10, 64); err != nil {
- return telegramCommandReply{Text: "TelegramID 必须是数字。"}
- }
- cfg := s.telegramChannelConfig(channel)
- ids := telegramConfiguredUserIDs(cfg["admin_user_ids"])
- seen := make(map[string]bool, len(ids)+1)
- var next []string
- for _, id := range ids {
- if id == tgID {
- seen[id] = true
- if add {
- next = append(next, id)
- }
- continue
- }
- if id != "" {
- next = append(next, id)
- }
- }
- if add && !seen[tgID] {
- next = append(next, tgID)
- }
- cfg["admin_user_ids"] = strings.Join(next, ",")
- raw, _ := json.Marshal(cfg)
- updated := *channel
- updated.Config = string(raw)
- if s.crypto != nil {
- updated.Config = s.crypto.Encrypt(updated.Config)
- }
- if err := s.repo.NotifyChannel.Update(ctx, &updated); err != nil {
- return telegramCommandReply{Text: "更新管理员列表失败:" + err.Error()}
- }
- if add {
- return telegramCommandReply{Text: "已添加 Bot 管理员:" + tgID + ""}
- }
- return telegramCommandReply{Text: "已移除 Bot 管理员:" + tgID + ""}
-}
-
-func (s *TelegramBotService) cmdMgoProtectedUser(ctx context.Context, args []string, protect bool) telegramCommandReply {
- if len(args) == 0 || strings.EqualFold(args[0], "list") {
- return s.cmdMgoProtectedUserList(ctx)
- }
- user := s.findMgoBotUser(ctx, args[0])
- if user == nil {
- return telegramCommandReply{Text: "未找到用户。"}
- }
- ids := ProtectedUserIDSet(ctx, s.repo)
- if protect {
- ids[user.ID] = struct{}{}
- if err := SaveProtectedUserIDSet(ctx, s.repo, ids); err != nil {
- return telegramCommandReply{Text: "保存保护名单失败:" + err.Error()}
- }
- return telegramCommandReply{Text: fmt.Sprintf("已加入保护名单:%s。\n该用户不会被 Bot 自动清理、批量禁用或删除。", user.Username)}
- }
- delete(ids, user.ID)
- if err := SaveProtectedUserIDSet(ctx, s.repo, ids); err != nil {
- return telegramCommandReply{Text: "保存保护名单失败:" + err.Error()}
- }
- return telegramCommandReply{Text: fmt.Sprintf("已移出保护名单:%s。", user.Username)}
-}
-
-func (s *TelegramBotService) cmdMgoProtectedUserList(ctx context.Context) telegramCommandReply {
- ids := ProtectedUserIDSet(ctx, s.repo)
- if len(ids) == 0 {
- return telegramCommandReply{Text: "保护名单为空。管理员和默认管理员始终自动保护。"}
- }
- names := make([]string, 0, len(ids))
- for id := range ids {
- if user, _ := s.repo.User.FindByID(ctx, id); user != nil {
- names = append(names, user.Username)
- } else {
- names = append(names, id+"(用户不存在)")
- }
- }
- sort.Strings(names)
- return telegramCommandReply{Text: fmt.Sprintf("保护名单:%d 个。\n%s", len(names), telegramInlineCodeList(names))}
-}
-
-func (s *TelegramBotService) cmdMgoBackupDB(ctx context.Context) telegramCommandReply {
- if s.backup == nil {
- return telegramCommandReply{Text: "备份服务暂不可用。"}
- }
- info, err := s.backup.Create(ctx)
- if err != nil {
- return telegramCommandReply{Text: "数据库备份失败:" + err.Error()}
- }
- return telegramCommandReply{Text: fmt.Sprintf("数据库备份完成:%s\n大小:%d bytes", info.Filename, info.Size)}
-}
-
-func (s *TelegramBotService) cmdMgoRestoreDB(ctx context.Context, args []string) telegramCommandReply {
- if s.backup == nil {
- return telegramCommandReply{Text: "备份服务暂不可用。"}
- }
- if len(args) == 0 || strings.EqualFold(args[0], "list") {
- items, err := s.backup.List()
- if err != nil {
- return telegramCommandReply{Text: "读取备份列表失败:" + err.Error()}
- }
- if len(items) == 0 {
- return telegramCommandReply{Text: "暂无数据库备份。可先使用 /backup_db 创建。"}
- }
- lines := make([]string, 0, minInt(len(items), 10))
- for i, item := range items {
- if i >= 10 {
- break
- }
- lines = append(lines, fmt.Sprintf("%s(%d bytes)", item.Filename, item.Size))
- }
- return telegramCommandReply{Text: "可恢复备份:\n" + telegramInlineCodeList(lines) + "\n\n恢复需要确认:/restore_from_db 文件名 confirm"}
- }
- if len(args) < 2 || !strings.EqualFold(args[len(args)-1], "confirm") {
- return telegramCommandReply{Text: "恢复数据库会覆盖当前数据,需要确认:/restore_from_db 文件名 confirm"}
- }
- filename := strings.TrimSpace(args[0])
- if err := s.backup.Restore(ctx, filename); err != nil {
- return telegramCommandReply{Text: "恢复失败:" + err.Error()}
- }
- return telegramCommandReply{Text: "数据库已从备份恢复,请重启 MediaStationGo 后生效。"}
-}
-
-func (s *TelegramBotService) cmdMgoSyncGroup(ctx context.Context, channel *model.NotifyChannel, args []string) telegramCommandReply {
- chatIDs := s.telegramMembershipChatIDs(channel)
- if len(chatIDs) == 0 {
- return telegramCommandReply{Text: "未配置可校验成员的群组/频道 ID。请在 Telegram 通知渠道设置 group_chat_id 或 channel_chat_id。"}
- }
- if strings.TrimSpace(s.telegramChannelConfig(channel)["bot_token"]) == "" {
- return telegramCommandReply{Text: "当前 Telegram 渠道未配置 bot_token,无法校验群成员。"}
- }
- var bindings []model.TelegramBinding
- if err := s.repo.DB.WithContext(ctx).Find(&bindings).Error; err != nil {
- return telegramCommandReply{Text: "读取绑定失败:" + err.Error()}
- }
- type staleBinding struct {
- User model.User
- Binding model.TelegramBinding
- }
- var stale []staleBinding
- for _, binding := range bindings {
- if binding.TelegramUserID == 0 || binding.UserID == "" {
- continue
- }
- user, _ := s.repo.User.FindByID(ctx, binding.UserID)
- if user == nil || UserIsProtectedAccount(ctx, s.repo, user) {
- continue
- }
- // 仅当所有绑定群组/频道都「查实不是成员」时才判定为可清理;
- // getChatMember 出错(membershipUnknown)时保守跳过,避免误删。
- confirmedNo := true
- for _, chatID := range chatIDs {
- if s.telegramChatMembership(ctx, channel, chatID, int(binding.TelegramUserID)) != membershipNo {
- confirmedNo = false
- break
- }
- }
- if confirmedNo {
- stale = append(stale, staleBinding{User: *user, Binding: binding})
- }
- }
- if len(stale) == 0 {
- return telegramCommandReply{Text: "所有已绑定账号都仍在配置的群组/频道中。"}
- }
- if len(args) >= 2 && strings.EqualFold(args[0], "delete") && strings.EqualFold(args[1], "confirm") {
- deleted := 0
- for _, item := range stale {
- _ = s.repo.UserDevice.DeleteByUser(ctx, item.User.ID)
- if err := s.repo.User.Delete(ctx, item.User.ID); err == nil {
- deleted++
- }
- }
- return telegramCommandReply{Text: fmt.Sprintf("已删除不在群组/频道中的普通账号:%d 个。", deleted)}
- }
- names := make([]string, 0, minInt(len(stale), 20))
- for i, item := range stale {
- if i >= 20 {
- break
- }
- names = append(names, fmt.Sprintf("%s(tg:%d)", item.User.Username, item.Binding.TelegramUserID))
- }
- return telegramCommandReply{Text: fmt.Sprintf("不在配置群组/频道中的绑定账号:%d 个。\n%s\n\n删除需确认:/syncgroupm delete confirm", len(stale), telegramInlineCodeList(names))}
-}
-
-func (s *TelegramBotService) telegramMembershipChatIDs(channel *model.NotifyChannel) []string {
- cfg := s.telegramChannelConfig(channel)
- seen := map[string]struct{}{}
- var out []string
- for _, key := range []string{"group_chat_id", "channel_chat_id", "command_chat_id"} {
- value := strings.TrimSpace(cfg[key])
- if value == "" {
- continue
- }
- if _, ok := seen[value]; !ok {
- seen[value] = struct{}{}
- out = append(out, value)
- }
- }
- if len(out) == 0 {
- if value := strings.TrimSpace(cfg["chat_id"]); strings.HasPrefix(value, "-") {
- out = append(out, value)
- }
- }
- return out
-}
-
func (s *TelegramBotService) cmdMgoUnsupported(name, replacement string) telegramCommandReply {
text := fmt.Sprintf("%s 已识别,但当前 Telegram Bot API 无法完整复刻该行为。", name)
if replacement != "" {
diff --git a/internal/service/telegram_mgo_membership.go b/internal/service/telegram_mgo_membership.go
new file mode 100644
index 0000000..0804483
--- /dev/null
+++ b/internal/service/telegram_mgo_membership.go
@@ -0,0 +1,92 @@
+package service
+
+import (
+ "context"
+ "fmt"
+ "strings"
+
+ "github.com/ShukeBta/MediaStationGo/internal/model"
+)
+
+func (s *TelegramBotService) cmdMgoSyncGroup(ctx context.Context, channel *model.NotifyChannel, args []string) telegramCommandReply {
+ chatIDs := s.telegramMembershipChatIDs(channel)
+ if len(chatIDs) == 0 {
+ return telegramCommandReply{Text: "未配置可校验成员的群组/频道 ID。请在 Telegram 通知渠道设置 group_chat_id 或 channel_chat_id。"}
+ }
+ if strings.TrimSpace(s.telegramChannelConfig(channel)["bot_token"]) == "" {
+ return telegramCommandReply{Text: "当前 Telegram 渠道未配置 bot_token,无法校验群成员。"}
+ }
+ var bindings []model.TelegramBinding
+ if err := s.repo.DB.WithContext(ctx).Find(&bindings).Error; err != nil {
+ return telegramCommandReply{Text: "读取绑定失败:" + err.Error()}
+ }
+ type staleBinding struct {
+ User model.User
+ Binding model.TelegramBinding
+ }
+ var stale []staleBinding
+ for _, binding := range bindings {
+ if binding.TelegramUserID == 0 || binding.UserID == "" {
+ continue
+ }
+ user, _ := s.repo.User.FindByID(ctx, binding.UserID)
+ if user == nil || UserIsProtectedAccount(ctx, s.repo, user) {
+ continue
+ }
+ // 仅当所有绑定群组/频道都「查实不是成员」时才判定为可清理;
+ // getChatMember 出错(membershipUnknown)时保守跳过,避免误删。
+ confirmedNo := true
+ for _, chatID := range chatIDs {
+ if s.telegramChatMembership(ctx, channel, chatID, int(binding.TelegramUserID)) != membershipNo {
+ confirmedNo = false
+ break
+ }
+ }
+ if confirmedNo {
+ stale = append(stale, staleBinding{User: *user, Binding: binding})
+ }
+ }
+ if len(stale) == 0 {
+ return telegramCommandReply{Text: "所有已绑定账号都仍在配置的群组/频道中。"}
+ }
+ if len(args) >= 2 && strings.EqualFold(args[0], "delete") && strings.EqualFold(args[1], "confirm") {
+ deleted := 0
+ for _, item := range stale {
+ _ = s.repo.UserDevice.DeleteByUser(ctx, item.User.ID)
+ if err := s.repo.User.Delete(ctx, item.User.ID); err == nil {
+ deleted++
+ }
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("已删除不在群组/频道中的普通账号:%d 个。", deleted)}
+ }
+ names := make([]string, 0, minInt(len(stale), 20))
+ for i, item := range stale {
+ if i >= 20 {
+ break
+ }
+ names = append(names, fmt.Sprintf("%s(tg:%d)", item.User.Username, item.Binding.TelegramUserID))
+ }
+ return telegramCommandReply{Text: fmt.Sprintf("不在配置群组/频道中的绑定账号:%d 个。\n%s\n\n删除需确认:/syncgroupm delete confirm", len(stale), telegramInlineCodeList(names))}
+}
+
+func (s *TelegramBotService) telegramMembershipChatIDs(channel *model.NotifyChannel) []string {
+ cfg := s.telegramChannelConfig(channel)
+ seen := map[string]struct{}{}
+ var out []string
+ for _, key := range []string{"group_chat_id", "channel_chat_id", "command_chat_id"} {
+ value := strings.TrimSpace(cfg[key])
+ if value == "" {
+ continue
+ }
+ if _, ok := seen[value]; !ok {
+ seen[value] = struct{}{}
+ out = append(out, value)
+ }
+ }
+ if len(out) == 0 {
+ if value := strings.TrimSpace(cfg["chat_id"]); strings.HasPrefix(value, "-") {
+ out = append(out, value)
+ }
+ }
+ return out
+}