From b9b7b7052ad87d293650d96f57fb38c78bc1ecc7 Mon Sep 17 00:00:00 2001 From: soldosluka857 Date: Sat, 30 May 2026 03:56:15 +0000 Subject: [PATCH] feat(organize/scan): transfer modes, seeding-safe relocation, inode dedup, incremental scanning - Organizer: add move/copy/hardlink/symlink transfer modes (default move), per-request target_path/transfer_mode overrides, and honor organize.target_dir / organize.transfer_mode settings. - keep_seeding (default on): escalate move->hardlink (cross-device->copy) so the qBittorrent source stays in place and continues seeding after organize. - qBittorrent SetLocation + POST /downloads/relocate to migrate whole torrents while keeping them seeding. - Scanner: FileID (device:inode) hardlink dedup to avoid duplicate recognition and double-counted storage; extract single-file ingest. - Watcher: recursive watch + incremental per-file ingest/remove instead of full re-scan; periodic full library scan now gated behind scan.periodic_enabled (default off) to reduce disk wear. - Telegram bot: handle callback_query in polling, declare allowed_updates in webhook, answer callbacks. - Frontend: settings for transfer mode / keep_seeding / periodic scan; organize panel target dir + transfer mode overrides. - Tests for transfer modes, organizer resolution, SetLocation, inode dedup, incremental ingest/remove. Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- internal/handler/downloads.go | 22 ++ internal/handler/handler.go | 1 + internal/handler/organizer.go | 25 ++- internal/model/model.go | 6 + internal/repository/repository.go | 4 + internal/service/downloads.go | 13 ++ internal/service/fileid_other.go | 9 + internal/service/fileid_unix.go | 30 +++ internal/service/organizer.go | 100 ++++++++- internal/service/organizer_transfer_test.go | 146 +++++++++++++ internal/service/qbittorrent.go | 48 +++++ internal/service/qbittorrent_test.go | 63 ++++++ internal/service/scanner.go | 216 +++++++++++++------ internal/service/scanner_incremental_test.go | 135 ++++++++++++ internal/service/scheduler.go | 55 +++-- internal/service/telegram_bot.go | 32 ++- internal/service/telegram_bot_user_test.go | 72 +++++++ internal/service/transfer.go | 98 +++++++++ internal/service/transfer_test.go | 123 +++++++++++ internal/service/watcher.go | 120 +++++++++-- web/src/api/tools.ts | 18 +- web/src/pages/SettingsPage.tsx | 24 +++ web/src/pages/ToolsPage.tsx | 43 +++- 23 files changed, 1284 insertions(+), 119 deletions(-) create mode 100644 internal/service/fileid_other.go create mode 100644 internal/service/fileid_unix.go create mode 100644 internal/service/organizer_transfer_test.go create mode 100644 internal/service/scanner_incremental_test.go create mode 100644 internal/service/transfer.go create mode 100644 internal/service/transfer_test.go diff --git a/internal/handler/downloads.go b/internal/handler/downloads.go index 29fa0f5..a1a8ff9 100644 --- a/internal/handler/downloads.go +++ b/internal/handler/downloads.go @@ -190,6 +190,28 @@ func deleteDownloadHandler(svc *service.Container) gin.HandlerFunc { } } +type relocateDownloadReq struct { + Hash string `json:"hash" binding:"required"` + Location string `json:"location" binding:"required"` +} + +// relocateDownloadHandler moves a torrent's data to a new directory while +// keeping it seeding (qBittorrent setLocation). +func relocateDownloadHandler(svc *service.Container) gin.HandlerFunc { + return func(c *gin.Context) { + var req relocateDownloadReq + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + if err := svc.Downloads.RelocateTorrent(c.Request.Context(), req.Hash, req.Location); err != nil { + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"hash": strings.TrimSpace(req.Hash), "location": strings.TrimSpace(req.Location)}) + } +} + func reloadDownloadConfigHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { if err := svc.Downloads.ReloadConfig(c.Request.Context()); err != nil { diff --git a/internal/handler/handler.go b/internal/handler/handler.go index 5e0befd..45c072b 100644 --- a/internal/handler/handler.go +++ b/internal/handler/handler.go @@ -109,6 +109,7 @@ func Register(r *gin.Engine, cfg *config.Config, log *zap.Logger, svc *service.C authed.GET("/downloads", requirePermission(svc, "can_manage_downloads"), listDownloadsHandler(svc)) authed.POST("/downloads", requirePermission(svc, "can_manage_downloads"), addDownloadHandler(svc)) authed.DELETE("/downloads/:hash", requirePermission(svc, "can_manage_downloads"), deleteDownloadHandler(svc)) + authed.POST("/downloads/relocate", requirePermission(svc, "can_manage_downloads"), relocateDownloadHandler(svc)) authed.POST("/downloads/reload", requirePermission(svc, "can_manage_downloads"), reloadDownloadConfigHandler(svc)) // Subscriptions. diff --git a/internal/handler/organizer.go b/internal/handler/organizer.go index bd48a2b..c7fe91f 100644 --- a/internal/handler/organizer.go +++ b/internal/handler/organizer.go @@ -3,15 +3,35 @@ package handler import ( "net/http" + "strings" "github.com/gin-gonic/gin" "github.com/ShukeBta/MediaStationGo/internal/service" ) +// organizeReq carries optional per-request overrides. 留空则沿用系统设置。 +type organizeReq struct { + TargetPath string `json:"target_path"` + TransferMode string `json:"transfer_mode"` +} + +// bindOrganizeOptions parses the optional JSON body into OrganizeOptions. +// A missing/empty body is fine — it means "use the configured defaults". +func bindOrganizeOptions(c *gin.Context) service.OrganizeOptions { + var req organizeReq + _ = c.ShouldBindJSON(&req) + opts := service.OrganizeOptions{TargetPath: strings.TrimSpace(req.TargetPath)} + if m := strings.TrimSpace(req.TransferMode); m != "" { + opts.TransferMode = service.TransferMode(m) + } + return opts +} + func organizeMediaHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { - dst, err := svc.Organizer.OrganizeMedia(c.Request.Context(), c.Param("id")) + opts := bindOrganizeOptions(c) + dst, err := svc.Organizer.OrganizeMediaWithOptions(c.Request.Context(), c.Param("id"), opts) if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return @@ -22,7 +42,8 @@ func organizeMediaHandler(svc *service.Container) gin.HandlerFunc { func organizeLibraryHandler(svc *service.Container) gin.HandlerFunc { return func(c *gin.Context) { - res, err := svc.Organizer.OrganizeLibrary(c.Request.Context(), c.Param("id")) + opts := bindOrganizeOptions(c) + res, err := svc.Organizer.OrganizeLibraryWithOptions(c.Request.Context(), c.Param("id"), opts) if err != nil { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return diff --git a/internal/model/model.go b/internal/model/model.go index 73f3c0b..3bd3d66 100644 --- a/internal/model/model.go +++ b/internal/model/model.go @@ -96,6 +96,12 @@ type Media struct { // Computed on-demand by the duplicate finder; format: "-". FileHash string `gorm:"index;size:64" json:"file_hash,omitempty"` + // FileID is a "device:inode" identity for the underlying file. Hardlinks + // to the same data share a FileID, letting the scanner skip re-importing a + // seeding source kept by keep_seeding and its organized hardlink as two + // separate items (avoids duplicate rows + double-counted storage). + FileID string `gorm:"index;size:64" json:"file_id,omitempty"` + // IsDuplicate flags this media as a duplicate of another media row. IsDuplicate bool `gorm:"default:false" json:"is_duplicate"` DuplicateOf string `gorm:"size:36" json:"duplicate_of,omitempty"` diff --git a/internal/repository/repository.go b/internal/repository/repository.go index e0f21e4..2db4413 100644 --- a/internal/repository/repository.go +++ b/internal/repository/repository.go @@ -307,6 +307,10 @@ func (r *MediaRepository) Upsert(ctx context.Context, m *model.Media) error { "container": m.Container, "deleted_at": nil, } + // 回填硬链接身份标识,便于后续扫描去重(避免重复识别/多倍占用)。 + if m.FileID != "" && m.FileID != existing.FileID { + updates["file_id"] = m.FileID + } if m.Title != "" { // scanner 给出的标题只是从路径推导,刮削后 title 已被替换为 // 真实剧名。仅在 existing 还停留在 'pending'/'' 时回填扫描标题, diff --git a/internal/service/downloads.go b/internal/service/downloads.go index 1dce4d2..88f9fbd 100644 --- a/internal/service/downloads.go +++ b/internal/service/downloads.go @@ -450,6 +450,19 @@ func (d *DownloadService) Delete(ctx context.Context, hash string, withFiles boo return d.qb.Delete(ctx, hash, withFiles) } +// RelocateTorrent moves a torrent's data to a new save directory while keeping +// it seeding (qBittorrent performs the physical move and resumes seeding). +// 用于「移动 PT 种子文件且转移后继续做种上传」的整盘迁移场景。 +func (d *DownloadService) RelocateTorrent(ctx context.Context, hash, location string) error { + if strings.TrimSpace(hash) == "" { + return errors.New("hash is required") + } + if strings.TrimSpace(location) == "" { + return errors.New("location is required") + } + return d.qb.SetLocation(ctx, hash, strings.TrimSpace(location)) +} + // poll fans out qBittorrent /torrents/info every 5 s as WS events. The // payload is opaque to the client; the React store merges by hash. func (d *DownloadService) poll(ctx context.Context) { diff --git a/internal/service/fileid_other.go b/internal/service/fileid_other.go new file mode 100644 index 0000000..f924880 --- /dev/null +++ b/internal/service/fileid_other.go @@ -0,0 +1,9 @@ +//go:build !unix + +package service + +// fileIdentity is a no-op on platforms without inode semantics; dedup by +// hardlink identity is simply disabled there. +func fileIdentity(path string) (string, bool) { + return "", false +} diff --git a/internal/service/fileid_unix.go b/internal/service/fileid_unix.go new file mode 100644 index 0000000..0dc0a7c --- /dev/null +++ b/internal/service/fileid_unix.go @@ -0,0 +1,30 @@ +//go:build unix + +package service + +import ( + "fmt" + "os" + "syscall" +) + +// fileIdentity returns a stable "device:inode" identifier for the file at +// path. Hardlinks to the same data share an identity, which lets the scanner +// avoid importing the same physical file twice (e.g. a seeding source kept by +// keep_seeding and its organized hardlink). ok is false when the identity +// cannot be determined. +func fileIdentity(path string) (string, bool) { + fi, err := os.Stat(path) + if err != nil { + return "", false + } + st, ok := fi.Sys().(*syscall.Stat_t) + if !ok || st == nil { + return "", false + } + // 单链接文件没有去重意义,避免给独立文件也打上可碰撞的标识。 + if st.Nlink < 2 { + return "", false + } + return fmt.Sprintf("%d:%d", uint64(st.Dev), uint64(st.Ino)), true +} diff --git a/internal/service/organizer.go b/internal/service/organizer.go index 6db105e..adc509b 100644 --- a/internal/service/organizer.go +++ b/internal/service/organizer.go @@ -48,11 +48,26 @@ type OrganizeResult struct { Errors []string `json:"errors,omitempty"` } +// OrganizeOptions carries per-request overrides for an organize operation. +// 空值表示沿用系统设置中的默认值。 +type OrganizeOptions struct { + // TargetPath 本次整理的目标根路径,覆盖 organize.target_dir 设置与媒体库路径。 + TargetPath string + // TransferMode 本次整理的转移方式,覆盖 organize.transfer_mode 设置。 + TransferMode TransferMode +} + // OrganizeMedia moves a single media file into the target library directory. // It auto-detects whether the media is a movie or TV episode based on the // parsed season/episode numbers and builds the destination path accordingly. // When smart classify is enabled, it adds a category subfolder (e.g., "华语电影"). func (o *OrganizerService) OrganizeMedia(ctx context.Context, mediaID string) (string, error) { + return o.OrganizeMediaWithOptions(ctx, mediaID, OrganizeOptions{}) +} + +// OrganizeMediaWithOptions is OrganizeMedia with per-request overrides for the +// target path and transfer mode. +func (o *OrganizerService) OrganizeMediaWithOptions(ctx context.Context, mediaID string, opts OrganizeOptions) (string, error) { m, err := o.repo.Media.FindByID(ctx, mediaID) if err != nil || m == nil { return "", errors.New("media not found") @@ -61,6 +76,8 @@ func (o *OrganizerService) OrganizeMedia(ctx context.Context, mediaID string) (s if err != nil || lib == nil { return "", errors.New("library not found") } + baseRoot := o.resolveBaseRoot(ctx, lib, opts.TargetPath) + mode := o.resolveTransferMode(ctx, opts.TransferMode) if isSeriesLibraryType(lib.Type) { if err := o.refreshEpisodeIdentity(m, lib); err != nil { return "", err @@ -77,19 +94,19 @@ func (o *OrganizerService) OrganizeMedia(ctx context.Context, mediaID string) (s var dst string if isSeriesLibraryType(lib.Type) { - // TV: {lib.Path}/[分类]/{Title}/Season XX/{Title} - SxxExx.ext + // TV: {baseRoot}/[分类]/{Title}/Season XX/{Title} - SxxExx.ext season := fmt.Sprintf("Season %02d", m.SeasonNum) epTag := fmt.Sprintf("S%02dE%02d", m.SeasonNum, m.EpisodeNum) - root := o.organizeRoot(lib.Path, lib.Type, category) + root := o.organizeRoot(baseRoot, lib.Type, category) dir := filepath.Join(categoryRoot(root, category), title, season) dst = filepath.Join(dir, fmt.Sprintf("%s - %s%s", title, epTag, ext)) } else { - // Movie: {lib.Path}/[分类]/{Title} ({Year})/{Title} ({Year}).ext + // Movie: {baseRoot}/[分类]/{Title} ({Year})/{Title} ({Year}).ext folder := title if m.Year > 0 { folder = fmt.Sprintf("%s (%d)", title, m.Year) } - root := o.organizeRoot(lib.Path, lib.Type, category) + root := o.organizeRoot(baseRoot, lib.Type, category) dir := filepath.Join(categoryRoot(root, category), folder) dst = filepath.Join(dir, folder+ext) } @@ -115,8 +132,9 @@ func (o *OrganizerService) OrganizeMedia(ctx context.Context, mediaID string) (s return "", err } - // Move the file (same filesystem = rename; cross-device = copy+delete). - if err := moveFile(m.Path, dst); err != nil { + // Transfer the file according to the resolved mode. move 删除源; + // copy/hardlink/symlink 保留源文件,从而让下载器可继续做种。 + if err := transferFile(m.Path, dst, mode); err != nil { return "", err } @@ -131,7 +149,7 @@ func (o *OrganizerService) OrganizeMedia(ctx context.Context, mediaID string) (s }).Error; err != nil { return dst, err } - if err := moveSidecarNFO(m.Path, dst); err != nil { + if err := transferSidecarNFO(m.Path, dst, mode); err != nil { o.log.Warn("organize sidecar nfo failed", zap.String("media", m.ID), zap.String("from", nfoPath(m.Path)), @@ -143,13 +161,69 @@ func (o *OrganizerService) OrganizeMedia(ctx context.Context, mediaID string) (s zap.String("from", m.Path), zap.String("to", dst), zap.String("category", category), + zap.String("mode", string(mode)), ) return dst, nil } +// resolveBaseRoot picks the organize target root: a per-request override +// wins, then the organize.target_dir setting, then the library's own path. +func (o *OrganizerService) resolveBaseRoot(ctx context.Context, lib *model.Library, override string) string { + if r := strings.TrimSpace(override); r != "" { + return r + } + if o.repo != nil && o.repo.Setting != nil { + if v, err := o.repo.Setting.Get(ctx, "organize.target_dir"); err == nil && strings.TrimSpace(v) != "" { + return strings.TrimSpace(v) + } + } + return lib.Path +} + +// resolveTransferMode picks the transfer mode: a per-request override wins, +// otherwise the organize.transfer_mode setting (default move). When the +// effective mode is move and 做种保种 (organize.keep_seeding) is enabled, it is +// upgraded to hardlink so the source stays in place for the torrent client. +func (o *OrganizerService) resolveTransferMode(ctx context.Context, override TransferMode) TransferMode { + mode := override + if mode == "" { + mode = TransferMove + if o.repo != nil && o.repo.Setting != nil { + if v, err := o.repo.Setting.Get(ctx, "organize.transfer_mode"); err == nil && strings.TrimSpace(v) != "" { + mode = parseTransferMode(v) + } + } + } + if mode == TransferMove && o.keepSeedingEnabled(ctx) { + // 移动会删除源文件导致 qBittorrent 停止做种;保种开启时改用硬链接 + //(跨盘自动退化为复制),既规范命名又保留源文件继续做种上传。 + return TransferHardlink + } + return mode +} + +// keepSeedingEnabled reports whether 做种保种 is on. Defaults to true so an +// unconfigured instance never silently breaks seeding on organize. +func (o *OrganizerService) keepSeedingEnabled(ctx context.Context) bool { + if o.repo == nil || o.repo.Setting == nil { + return true + } + v, err := o.repo.Setting.Get(ctx, "organize.keep_seeding") + if err != nil || strings.TrimSpace(v) == "" { + return true + } + return v == "true" || v == "1" || v == "on" +} + // OrganizeLibrary organizes every media row in a library whose file is // not already in the expected path structure. func (o *OrganizerService) OrganizeLibrary(ctx context.Context, libraryID string) (*OrganizeResult, error) { + return o.OrganizeLibraryWithOptions(ctx, libraryID, OrganizeOptions{}) +} + +// OrganizeLibraryWithOptions is OrganizeLibrary with per-request overrides for +// the target path and transfer mode. +func (o *OrganizerService) OrganizeLibraryWithOptions(ctx context.Context, libraryID string, opts OrganizeOptions) (*OrganizeResult, error) { lib, err := o.repo.Library.FindByID(ctx, libraryID) if err != nil || lib == nil { return nil, errors.New("library not found") @@ -160,13 +234,15 @@ func (o *OrganizerService) OrganizeLibrary(ctx context.Context, libraryID string Find(&rows).Error; err != nil { return nil, err } + // 已位于整理目标根下的文件视为已整理;目标根受 target_path 覆盖与设置影响。 + baseRoot := o.resolveBaseRoot(ctx, lib, opts.TargetPath) res := &OrganizeResult{} for i := range rows { - if pathWithin(rows[i].Path, lib.Path) { + if pathWithin(rows[i].Path, baseRoot) { res.Skipped++ continue } - dst, err := o.OrganizeMedia(ctx, rows[i].ID) + dst, err := o.OrganizeMediaWithOptions(ctx, rows[i].ID, opts) if err != nil { res.Errors = append(res.Errors, fmt.Sprintf("%s: %s", rows[i].Title, err.Error())) continue @@ -247,7 +323,9 @@ func moveFile(src, dst string) error { return os.Remove(src) } -func moveSidecarNFO(srcMedia, dstMedia string) error { +// transferSidecarNFO moves/copies/links the .nfo sidecar alongside its media +// using the same transfer mode, so metadata follows the organized file. +func transferSidecarNFO(srcMedia, dstMedia string, mode TransferMode) error { src := nfoPath(srcMedia) dst := nfoPath(dstMedia) if src == dst { @@ -265,7 +343,7 @@ func moveSidecarNFO(srcMedia, dstMedia string) error { if err := os.MkdirAll(filepath.Dir(dst), 0o755); err != nil { return err } - return moveFile(src, dst) + return transferFile(src, dst, mode) } // sanitizeFilename removes characters not safe for filesystem names. diff --git a/internal/service/organizer_transfer_test.go b/internal/service/organizer_transfer_test.go new file mode 100644 index 0000000..82c0857 --- /dev/null +++ b/internal/service/organizer_transfer_test.go @@ -0,0 +1,146 @@ +package service + +import ( + "os" + "path/filepath" + "strings" + "testing" + + "github.com/glebarez/sqlite" + "go.uber.org/zap" + "gorm.io/gorm" + + "github.com/ShukeBta/MediaStationGo/internal/config" + "github.com/ShukeBta/MediaStationGo/internal/model" + "github.com/ShukeBta/MediaStationGo/internal/repository" +) + +func newOrganizerTestRepo(t *testing.T) *repository.Container { + t.Helper() + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&model.Library{}, &model.Media{}, &model.Setting{}); err != nil { + t.Fatal(err) + } + return repository.New(db) +} + +func TestOrganizeMediaHonorsTargetDirAndCopyMode(t *testing.T) { + root := t.TempDir() + srcDir := filepath.Join(root, "downloads") + if err := os.MkdirAll(srcDir, 0o755); err != nil { + t.Fatal(err) + } + source := filepath.Join(srcDir, "Some Movie.mkv") + if err := os.WriteFile(source, []byte("movie"), 0o644); err != nil { + t.Fatal(err) + } + target := filepath.Join(root, "custom-library") + + repos := newOrganizerTestRepo(t) + if err := repos.Setting.Set(t.Context(), "organize.target_dir", target); err != nil { + t.Fatal(err) + } + lib := model.Library{Name: "Movies", Path: filepath.Join(root, "lib"), Type: "movie", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + media := model.Media{LibraryID: lib.ID, Title: "Some Movie", Path: source, Year: 2020, Container: "mkv", ScrapeStatus: "matched"} + if err := repos.Media.Upsert(t.Context(), &media); err != nil { + t.Fatal(err) + } + + org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + dst, err := org.OrganizeMediaWithOptions(t.Context(), media.ID, OrganizeOptions{TransferMode: TransferCopy}) + if err != nil { + t.Fatalf("organize: %v", err) + } + if !strings.HasPrefix(dst, target) { + t.Fatalf("dst %q should be under target dir %q", dst, target) + } + if _, err := os.Stat(source); err != nil { + t.Fatalf("copy mode must keep source: %v", err) + } + if _, err := os.Stat(dst); err != nil { + t.Fatalf("organized file missing: %v", err) + } +} + +func TestOrganizeMediaKeepSeedingUpgradesMoveToHardlink(t *testing.T) { + root := t.TempDir() + srcDir := filepath.Join(root, "downloads") + if err := os.MkdirAll(srcDir, 0o755); err != nil { + t.Fatal(err) + } + source := filepath.Join(srcDir, "Some Movie.mkv") + if err := os.WriteFile(source, []byte("movie"), 0o644); err != nil { + t.Fatal(err) + } + + repos := newOrganizerTestRepo(t) + // 默认转移方式为 move,做种保种默认开启 → 应升级为硬链接(保留源文件)。 + lib := model.Library{Name: "Movies", Path: filepath.Join(root, "lib"), Type: "movie", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + media := model.Media{LibraryID: lib.ID, Title: "Some Movie", Path: source, Year: 2021, Container: "mkv", ScrapeStatus: "matched"} + if err := repos.Media.Upsert(t.Context(), &media); err != nil { + t.Fatal(err) + } + + org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + dst, err := org.OrganizeMedia(t.Context(), media.ID) + if err != nil { + t.Fatalf("organize: %v", err) + } + si, err := os.Stat(source) + if err != nil { + t.Fatalf("keep_seeding default should keep source for seeding: %v", err) + } + di, err := os.Stat(dst) + if err != nil { + t.Fatalf("organized file missing: %v", err) + } + if !os.SameFile(si, di) { + t.Fatal("organized file should be a hardlink of the source (same inode)") + } +} + +func TestOrganizeMediaKeepSeedingOffPerformsTrueMove(t *testing.T) { + root := t.TempDir() + srcDir := filepath.Join(root, "downloads") + if err := os.MkdirAll(srcDir, 0o755); err != nil { + t.Fatal(err) + } + source := filepath.Join(srcDir, "Some Movie.mkv") + if err := os.WriteFile(source, []byte("movie"), 0o644); err != nil { + t.Fatal(err) + } + + repos := newOrganizerTestRepo(t) + if err := repos.Setting.Set(t.Context(), "organize.keep_seeding", "false"); err != nil { + t.Fatal(err) + } + lib := model.Library{Name: "Movies", Path: filepath.Join(root, "lib"), Type: "movie", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + media := model.Media{LibraryID: lib.ID, Title: "Some Movie", Path: source, Year: 2022, Container: "mkv", ScrapeStatus: "matched"} + if err := repos.Media.Upsert(t.Context(), &media); err != nil { + t.Fatal(err) + } + + org := NewOrganizerService(&config.Config{}, zap.NewNop(), repos) + dst, err := org.OrganizeMedia(t.Context(), media.ID) + if err != nil { + t.Fatalf("organize: %v", err) + } + if _, err := os.Stat(source); !os.IsNotExist(err) { + t.Fatalf("keep_seeding off should remove source, stat err = %v", err) + } + if _, err := os.Stat(dst); err != nil { + t.Fatalf("organized file missing: %v", err) + } +} diff --git a/internal/service/qbittorrent.go b/internal/service/qbittorrent.go index 25e54a4..8d07ca6 100644 --- a/internal/service/qbittorrent.go +++ b/internal/service/qbittorrent.go @@ -454,6 +454,54 @@ func (q *QBitClient) Delete(ctx context.Context, hash string, deleteFiles bool) return nil } +// SetLocation moves a torrent's data to a new save directory via +// POST /api/v2/torrents/setLocation. qBittorrent performs the physical move +// itself and keeps seeding from the new location — this is the seeding-safe +// way to relocate downloaded PT files. location must be an absolute path the +// qBittorrent process can write to. +func (q *QBitClient) SetLocation(ctx context.Context, hash, location string) error { + if strings.TrimSpace(hash) == "" { + return errors.New("qbittorrent setLocation: empty hash") + } + if strings.TrimSpace(location) == "" { + return errors.New("qbittorrent setLocation: empty location") + } + q.mu.Lock() + defer q.mu.Unlock() + if err := q.ensureAuth(ctx); err != nil { + return err + } + form := url.Values{} + form.Set("hashes", hash) + form.Set("location", location) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, + strings.TrimRight(q.cfg.BaseURL, "/")+"/api/v2/torrents/setLocation", + strings.NewReader(form.Encode()), + ) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/x-www-form-urlencoded") + req.Header.Set("Referer", q.cfg.BaseURL) + resp, err := q.client.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + body, _ := io.ReadAll(resp.Body) + if resp.StatusCode == http.StatusBadRequest { + return errors.New("qbittorrent setLocation: 保存路径无效") + } + if resp.StatusCode == http.StatusConflict { + return errors.New("qbittorrent setLocation: 无法写入目标路径 (权限或磁盘问题)") + } + if resp.StatusCode >= 400 { + return fmt.Errorf("qbittorrent setLocation: HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) + } + q.log.Info("qbittorrent: torrent relocated", zap.String("hash", hash), zap.String("location", location)) + return nil +} + // ensureAuth makes sure we have a valid SID cookie. Cheap on the happy // path; logs in transparently otherwise. func (q *QBitClient) ensureAuth(ctx context.Context) error { diff --git a/internal/service/qbittorrent_test.go b/internal/service/qbittorrent_test.go index b4124d0..c2ab2db 100644 --- a/internal/service/qbittorrent_test.go +++ b/internal/service/qbittorrent_test.go @@ -212,3 +212,66 @@ func multipartHasTorrentFile(reader *multipart.Reader) bool { } } } + +func TestQBitSetLocationPostsHashAndLocation(t *testing.T) { + var gotHash, gotLocation string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/api/v2/auth/login": + _, _ = w.Write([]byte("Ok.")) + case "/api/v2/torrents/setLocation": + if err := r.ParseForm(); err != nil { + http.Error(w, "bad form", http.StatusBadRequest) + return + } + gotHash = r.PostFormValue("hashes") + gotLocation = r.PostFormValue("location") + _, _ = w.Write([]byte("Ok.")) + default: + http.NotFound(w, r) + } + })) + defer server.Close() + + client := NewQBitClient(zap.NewNop(), QBitConfig{ + BaseURL: server.URL, + Username: "admin", + Password: "adminadmin", + }) + if err := client.SetLocation(context.Background(), "abc123", "/data/media/Movie"); err != nil { + t.Fatalf("setLocation: %v", err) + } + if gotHash != "abc123" { + t.Fatalf("hashes = %q, want abc123", gotHash) + } + if gotLocation != "/data/media/Movie" { + t.Fatalf("location = %q, want /data/media/Movie", gotLocation) + } +} + +func TestQBitSetLocationSurfacesConflict(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/api/v2/auth/login": + _, _ = w.Write([]byte("Ok.")) + case "/api/v2/torrents/setLocation": + http.Error(w, "cannot write", http.StatusConflict) + default: + http.NotFound(w, r) + } + })) + defer server.Close() + + client := NewQBitClient(zap.NewNop(), QBitConfig{ + BaseURL: server.URL, + Username: "admin", + Password: "adminadmin", + }) + err := client.SetLocation(context.Background(), "abc123", "/data/media/Movie") + if err == nil { + t.Fatal("expected error on 409 conflict") + } + if !strings.Contains(err.Error(), "无法写入目标路径") { + t.Fatalf("unexpected error: %v", err) + } +} diff --git a/internal/service/scanner.go b/internal/service/scanner.go index 903ad33..3842004 100644 --- a/internal/service/scanner.go +++ b/internal/service/scanner.go @@ -72,6 +72,7 @@ type ScanResult struct { Visited int `json:"visited"` Added int `json:"added"` Updated int `json:"updated"` + Skipped int `json:"skipped"` Probed int `json:"probed"` LocalMetadata int `json:"local_metadata"` Removed int64 `json:"removed"` @@ -85,6 +86,7 @@ func (s *ScannerService) ScanLibrary(ctx context.Context, libraryID string) (*Sc } res := &ScanResult{LibraryID: lib.ID} seen := make(map[string]struct{}) + seenInodes := make(map[string]string) walkFn := func(path string, info walkInfo) error { if info.isDir { @@ -94,70 +96,8 @@ func (s *ScannerService) ScanLibrary(ctx context.Context, libraryID string) (*Sc if _, ok := videoExtensions[ext]; !ok { return nil } - res.Visited++ seen[filepath.Clean(path)] = struct{}{} - isNewMedia := !s.mediaPathExists(ctx, path) - - title, year := CleanQuery(path) - if title == "" { - title = strings.TrimSuffix(filepath.Base(path), ext) - } - - m := &model.Media{ - LibraryID: lib.ID, - Title: title, - Year: year, - Path: path, - SizeBytes: info.size, - Container: strings.TrimPrefix(ext, "."), - } - - parsedSeason, parsedEpisode := ParseEpisode(path) - m.SeasonNum = parsedSeason - m.EpisodeNum = parsedEpisode - - if local, err := ReadLocalMetadata(path, lib.Path, librarySupportsSeasons(lib) || parsedSeason > 0 || parsedEpisode > 0); err == nil && local != nil { - applyLocalMetadata(m, local) - res.LocalMetadata++ - } else if err != nil { - s.log.Warn("read local metadata failed", zap.String("path", path), zap.Error(err)) - } - - // Best-effort ffprobe; failure does not abort the file. - if s.probe != nil { - if probe, err := s.probe.Probe(ctx, path); err == nil && probe != nil { - m.DurationSec = probe.DurationSec - m.Width = probe.Width - m.Height = probe.Height - m.VideoCodec = probe.VideoCodec - m.AudioCodec = probe.AudioCodec - if probe.Container != "" { - m.Container = probe.Container - } - res.Probed++ - } else if err != nil { - s.log.Debug("ffprobe failed", zap.String("path", path), zap.Error(err)) - } - } - - if err := s.repo.Media.Upsert(ctx, m); err != nil { - s.log.Warn("upsert media failed", zap.String("path", path), zap.Error(err)) - return nil - } - if isNewMedia { - res.Added++ - } else { - res.Updated++ - } - s.hub.Publish("scan", map[string]any{ - "library_id": lib.ID, - "path": path, - "visited": res.Visited, - "added": res.Added, - "updated": res.Updated, - "probed": res.Probed, - "local_meta": res.LocalMetadata, - }) + s.ingestFile(ctx, lib, path, info.size, seenInodes, res) return nil } @@ -194,6 +134,156 @@ func (s *ScannerService) ScanLibrary(ctx context.Context, libraryID string) (*Sc return res, nil } +// IngestPath ingests a single file into the given library without walking the +// whole tree. Used by the watcher for incremental, event-driven additions so +// adding one new file no longer triggers a full library re-scan (减少硬盘损耗). +// Non-video files and directories are ignored. Returns true if a media row was +// added or updated. +func (s *ScannerService) IngestPath(ctx context.Context, libraryID, path string) (bool, error) { + lib, err := s.repo.Library.FindByID(ctx, libraryID) + if err != nil || lib == nil { + return false, err + } + fi, err := os.Stat(path) + if err != nil || fi.IsDir() { + return false, err + } + ext := strings.ToLower(filepath.Ext(path)) + if _, ok := videoExtensions[ext]; !ok { + return false, nil + } + res := &ScanResult{LibraryID: lib.ID} + s.ingestFile(ctx, lib, path, fi.Size(), make(map[string]string), res) + return res.Added+res.Updated > 0, nil +} + +// RemovePath deletes the media row for a path that has disappeared from disk +// (incremental delete used by the watcher on Remove/Rename events). +func (s *ScannerService) RemovePath(ctx context.Context, path string) (int64, error) { + if _, err := os.Stat(path); err == nil { + return 0, nil // still exists; nothing to remove + } + res := s.repo.DB.WithContext(ctx). + Where("path = ?", path). + Delete(&model.Media{}) + return res.RowsAffected, res.Error +} + +// ingestFile upserts a single media file. seenInodes dedups hardlinks within a +// single scan; pass a fresh map for one-off ingests. It mutates res counters. +func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, path string, size int64, seenInodes map[string]string, res *ScanResult) { + res.Visited++ + ext := strings.ToLower(filepath.Ext(path)) + + // Hardlink dedup: a seeding source kept by keep_seeding shares its inode + // with the organized hardlink. Importing both would create duplicate rows + // and double-count storage, so skip any file whose identity we've already + // taken (within this scan or via an existing DB row pointing elsewhere). + fileID, hasID := fileIdentity(path) + if hasID { + if first, ok := seenInodes[fileID]; ok && first != path { + res.Skipped++ + s.log.Debug("scan skip hardlink duplicate", + zap.String("path", path), zap.String("primary", first)) + return + } + if other, ok := s.duplicateByFileID(ctx, fileID, path); ok { + res.Skipped++ + s.log.Debug("scan skip hardlink duplicate (existing)", + zap.String("path", path), zap.String("primary", other)) + return + } + seenInodes[fileID] = path + } + + isNewMedia := !s.mediaPathExists(ctx, path) + + title, year := CleanQuery(path) + if title == "" { + title = strings.TrimSuffix(filepath.Base(path), ext) + } + + m := &model.Media{ + LibraryID: lib.ID, + Title: title, + Year: year, + Path: path, + SizeBytes: size, + Container: strings.TrimPrefix(ext, "."), + FileID: fileID, + } + + parsedSeason, parsedEpisode := ParseEpisode(path) + m.SeasonNum = parsedSeason + m.EpisodeNum = parsedEpisode + + if local, err := ReadLocalMetadata(path, lib.Path, librarySupportsSeasons(lib) || parsedSeason > 0 || parsedEpisode > 0); err == nil && local != nil { + applyLocalMetadata(m, local) + res.LocalMetadata++ + } else if err != nil { + s.log.Warn("read local metadata failed", zap.String("path", path), zap.Error(err)) + } + + // Best-effort ffprobe; failure does not abort the file. + if s.probe != nil { + if probe, err := s.probe.Probe(ctx, path); err == nil && probe != nil { + m.DurationSec = probe.DurationSec + m.Width = probe.Width + m.Height = probe.Height + m.VideoCodec = probe.VideoCodec + m.AudioCodec = probe.AudioCodec + if probe.Container != "" { + m.Container = probe.Container + } + res.Probed++ + } else if err != nil { + s.log.Debug("ffprobe failed", zap.String("path", path), zap.Error(err)) + } + } + + if err := s.repo.Media.Upsert(ctx, m); err != nil { + s.log.Warn("upsert media failed", zap.String("path", path), zap.Error(err)) + return + } + if isNewMedia { + res.Added++ + } else { + res.Updated++ + } + s.hub.Publish("scan", map[string]any{ + "library_id": lib.ID, + "path": path, + "visited": res.Visited, + "added": res.Added, + "updated": res.Updated, + "probed": res.Probed, + "local_meta": res.LocalMetadata, + }) +} + +// duplicateByFileID reports an existing media path that shares the given inode +// identity but lives at a different path and still exists on disk. +func (s *ScannerService) duplicateByFileID(ctx context.Context, fileID, path string) (string, bool) { + if fileID == "" { + return "", false + } + var rows []model.Media + if err := s.repo.DB.WithContext(ctx). + Where("file_id = ? AND path <> ?", fileID, path). + Limit(8).Find(&rows).Error; err != nil { + return "", false + } + for _, r := range rows { + if r.Path == "" { + continue + } + if _, err := os.Stat(r.Path); err == nil { + return r.Path, true + } + } + return "", false +} + func (s *ScannerService) mediaPathExists(ctx context.Context, path string) bool { var count int64 err := s.repo.DB.WithContext(ctx).Unscoped().Model(&model.Media{}). diff --git a/internal/service/scanner_incremental_test.go b/internal/service/scanner_incremental_test.go new file mode 100644 index 0000000..d1a180f --- /dev/null +++ b/internal/service/scanner_incremental_test.go @@ -0,0 +1,135 @@ +package service + +import ( + "os" + "path/filepath" + "testing" + + "github.com/glebarez/sqlite" + "go.uber.org/zap" + "gorm.io/gorm" + + "github.com/ShukeBta/MediaStationGo/internal/config" + "github.com/ShukeBta/MediaStationGo/internal/model" + "github.com/ShukeBta/MediaStationGo/internal/repository" +) + +func newScannerTestEnv(t *testing.T) (*ScannerService, *repository.Container) { + t.Helper() + db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&model.Library{}, &model.Media{}, &model.Setting{}); err != nil { + t.Fatal(err) + } + repos := repository.New(db) + sc := NewScannerService(&config.Config{}, zap.NewNop(), repos, NewHub(zap.NewNop()), nil, nil) + return sc, repos +} + +func countMedia(t *testing.T, repos *repository.Container) int64 { + t.Helper() + var n int64 + if err := repos.DB.Model(&model.Media{}).Count(&n).Error; err != nil { + t.Fatal(err) + } + return n +} + +func TestIngestPathAddsSingleFile(t *testing.T) { + sc, repos := newScannerTestEnv(t) + root := t.TempDir() + lib := model.Library{Name: "Movies", Path: root, Type: "movie", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + file := filepath.Join(root, "Some Movie (2021).mkv") + if err := os.WriteFile(file, []byte("data"), 0o644); err != nil { + t.Fatal(err) + } + added, err := sc.IngestPath(t.Context(), lib.ID, file) + if err != nil { + t.Fatalf("ingest: %v", err) + } + if !added { + t.Fatal("expected file to be added") + } + if got := countMedia(t, repos); got != 1 { + t.Fatalf("media count = %d, want 1", got) + } + // Non-video file is ignored. + other := filepath.Join(root, "notes.txt") + if err := os.WriteFile(other, []byte("x"), 0o644); err != nil { + t.Fatal(err) + } + if added, _ := sc.IngestPath(t.Context(), lib.ID, other); added { + t.Fatal("non-video file should not be ingested") + } +} + +func TestRemovePathDeletesVanishedMedia(t *testing.T) { + sc, repos := newScannerTestEnv(t) + root := t.TempDir() + lib := model.Library{Name: "Movies", Path: root, Type: "movie", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + file := filepath.Join(root, "Gone (2020).mkv") + if err := os.WriteFile(file, []byte("data"), 0o644); err != nil { + t.Fatal(err) + } + if _, err := sc.IngestPath(t.Context(), lib.ID, file); err != nil { + t.Fatal(err) + } + if countMedia(t, repos) != 1 { + t.Fatal("expected 1 media before removal") + } + // A still-present file is not removed. + if removed, _ := sc.RemovePath(t.Context(), file); removed != 0 { + t.Fatalf("present file should not be removed, got %d", removed) + } + if err := os.Remove(file); err != nil { + t.Fatal(err) + } + removed, err := sc.RemovePath(t.Context(), file) + if err != nil { + t.Fatalf("remove: %v", err) + } + if removed != 1 { + t.Fatalf("removed = %d, want 1", removed) + } + if countMedia(t, repos) != 0 { + t.Fatal("expected 0 media after removal") + } +} + +// TestScanSkipsHardlinkDuplicate verifies that a hardlink (same inode) kept +// for seeding is not imported as a second media item, preventing duplicate +// recognition and double-counted storage. +func TestScanSkipsHardlinkDuplicate(t *testing.T) { + sc, repos := newScannerTestEnv(t) + root := t.TempDir() + lib := model.Library{Name: "Movies", Path: root, Type: "movie", Enabled: true} + if err := repos.Library.Create(t.Context(), &lib); err != nil { + t.Fatal(err) + } + orig := filepath.Join(root, "Movie (2019).mkv") + if err := os.WriteFile(orig, []byte("payload"), 0o644); err != nil { + t.Fatal(err) + } + linked := filepath.Join(root, "Movie (2019) [organized].mkv") + if err := os.Link(orig, linked); err != nil { + t.Skipf("hardlinks unsupported on this fs: %v", err) + } + res, err := sc.ScanLibrary(t.Context(), lib.ID) + if err != nil { + t.Fatalf("scan: %v", err) + } + if got := countMedia(t, repos); got != 1 { + t.Fatalf("media count = %d, want 1 (hardlink should be deduped)", got) + } + if res.Skipped != 1 { + t.Fatalf("Skipped = %d, want 1", res.Skipped) + } +} diff --git a/internal/service/scheduler.go b/internal/service/scheduler.go index 2a2ed14..f612c70 100644 --- a/internal/service/scheduler.go +++ b/internal/service/scheduler.go @@ -3,20 +3,20 @@ // SchedulerService runs five recurring background jobs that keep the // library up-to-date without operator intervention: // -// library_scan every 60 min — re-scan every enabled library so -// newly-copied files are picked up. -// subscription_pull every 30 min — re-poll RSS feeds (in addition to -// the existing SubscriptionService -// internal timer). -// download_sync every 30 s — refresh the qBittorrent torrent -// list (already covered by the -// download poller, kept here as a -// watchdog). -// transcode_cleanup every 24 h — purge HLS transcode artefacts -// older than 24 h. -// recycle_purge every 24 h — empty the recycle bin of rows -// soft-deleted more than 30 days -// ago. +// library_scan every 60 min — re-scan every enabled library so +// newly-copied files are picked up. +// subscription_pull every 30 min — re-poll RSS feeds (in addition to +// the existing SubscriptionService +// internal timer). +// download_sync every 30 s — refresh the qBittorrent torrent +// list (already covered by the +// download poller, kept here as a +// watchdog). +// transcode_cleanup every 24 h — purge HLS transcode artefacts +// older than 24 h. +// recycle_purge every 24 h — empty the recycle bin of rows +// soft-deleted more than 30 days +// ago. // // Each job runs at most once at a time (an in-flight run blocks the // next tick). All work happens on a long-lived background context so @@ -25,6 +25,7 @@ package service import ( "context" + "strings" "sync" "time" @@ -191,7 +192,14 @@ func (s *SchedulerService) runOnce(ctx context.Context, j *scheduledJob) error { } // jobScanLibraries re-walks every enabled library. +// +// 默认关闭:文件变更由 WatcherService 增量入库,无需周期性全量重扫。 +// 仅当用户在设置中显式开启 scan.periodic_enabled 时才执行整库重扫, +// 避免对硬盘的高频反复读取造成损伤(用户明确要求)。 func (s *SchedulerService) jobScanLibraries(ctx context.Context) error { + if !s.periodicScanEnabled(ctx) { + return nil + } libs, err := s.repo.Library.List(ctx) if err != nil { return err @@ -208,6 +216,25 @@ func (s *SchedulerService) jobScanLibraries(ctx context.Context) error { return nil } +// periodicScanEnabled reports whether the operator opted into periodic full +// library re-scans. Defaults to false so the incremental watcher is the only +// thing touching the disk under normal operation. +func (s *SchedulerService) periodicScanEnabled(ctx context.Context) bool { + if s.repo == nil || s.repo.Setting == nil { + return false + } + v, err := s.repo.Setting.Get(ctx, "scan.periodic_enabled") + if err != nil { + return false + } + switch strings.ToLower(strings.TrimSpace(v)) { + case "1", "true", "yes", "on", "enabled": + return true + default: + return false + } +} + // jobCleanTranscodeCache deletes HLS artefacts older than 24h. func (s *SchedulerService) jobCleanTranscodeCache(ctx context.Context) error { if s.cacheDir == "" { diff --git a/internal/service/telegram_bot.go b/internal/service/telegram_bot.go index e8f0951..da92bf1 100644 --- a/internal/service/telegram_bot.go +++ b/internal/service/telegram_bot.go @@ -573,7 +573,7 @@ func (s *TelegramBotService) pollLoop(ctx context.Context, cfg map[string]string if upd.UpdateID >= int(offset) { offset = int64(upd.UpdateID) + 1 } - if upd.Message == nil || upd.Message.Text == "" { + if !telegramUpdateActionable(upd) { continue } go func(u TelegramUpdate) { @@ -584,6 +584,16 @@ func (s *TelegramBotService) pollLoop(ctx context.Context, cfg map[string]string } } +// telegramUpdateActionable 判断一条 update 是否需要分发处理。 +// 长轮询默认会返回 message 与 callback_query 两类更新;命令消息需有文本, +// 而内联按钮回调(callback_query)必须被分发,否则成人目录显隐开关会失效。 +func telegramUpdateActionable(upd TelegramUpdate) bool { + if upd.CallbackQuery != nil { + return true + } + return upd.Message != nil && upd.Message.Text != "" +} + func telegramPollingRequest(ctx context.Context, clients []*http.Client, pollURL, body string) ([]byte, error) { var lastErr error for _, client := range clients { @@ -714,6 +724,8 @@ func (s *TelegramBotService) handleCallback(ctx context.Context, cb *TelegramCal if channel == nil { channel = s.findChannelByChatID(ctx, cb.Message.Chat.ID) } + // 立即应答回调,关闭按钮上的加载状态,避免客户端长时间转圈。 + s.answerCallback(ctx, channel, cb.ID) switch strings.TrimSpace(cb.Data) { case "adult_toggle": reply := s.cmdHideAdult(ctx, &msg, nil) @@ -724,6 +736,22 @@ func (s *TelegramBotService) handleCallback(ctx context.Context, cb *TelegramCal return nil } +// answerCallback 应答 Telegram 回调查询,关闭按钮上的加载提示。 +func (s *TelegramBotService) answerCallback(ctx context.Context, channel *model.NotifyChannel, callbackID string) { + if channel == nil || strings.TrimSpace(callbackID) == "" { + return + } + cfg := s.telegramChannelConfig(channel) + if strings.TrimSpace(cfg["bot_token"]) == "" { + return + } + if err := telegramPostJSON(ctx, cfg, "answerCallbackQuery", map[string]interface{}{ + "callback_query_id": callbackID, + }, 8*time.Second); err != nil { + s.log.Debug("telegram answerCallbackQuery failed", zap.Error(sanitizeTelegramError(err))) + } +} + func (s *TelegramBotService) telegramBinding(ctx context.Context, telegramUserID int) *model.TelegramBinding { if telegramUserID == 0 { return nil @@ -938,7 +966,7 @@ func userNameOrFallback(user *model.User) string { func (s *TelegramBotService) SetWebhook(ctx context.Context, botToken, webhookURL string) error { payload, _ := json.Marshal(map[string]interface{}{ "url": webhookURL, - "allowed_updates": []string{"message"}, + "allowed_updates": []string{"message", "callback_query"}, }) cfg := map[string]string{"bot_token": botToken} apiURL, err := telegramMethodURL(cfg, botToken, "setWebhook") diff --git a/internal/service/telegram_bot_user_test.go b/internal/service/telegram_bot_user_test.go index f0dc1b9..78fab8d 100644 --- a/internal/service/telegram_bot_user_test.go +++ b/internal/service/telegram_bot_user_test.go @@ -1,6 +1,7 @@ package service import ( + "encoding/json" "strings" "testing" @@ -9,6 +10,77 @@ import ( "github.com/ShukeBta/MediaStationGo/internal/model" ) +func TestTelegramUpdateActionableDispatchesCallbackQuery(t *testing.T) { + if !telegramUpdateActionable(TelegramUpdate{CallbackQuery: &TelegramCallbackQuery{Data: "adult_toggle"}}) { + t.Fatal("callback_query update must be dispatched, otherwise inline buttons break") + } + if !telegramUpdateActionable(TelegramUpdate{Message: &TelegramMessage{Text: "/help"}}) { + t.Fatal("text command message must be dispatched") + } + if telegramUpdateActionable(TelegramUpdate{}) { + t.Fatal("empty update must be skipped") + } + if telegramUpdateActionable(TelegramUpdate{Message: &TelegramMessage{}}) { + t.Fatal("message without text must be skipped") + } +} + +func TestTelegramCallbackTogglesAdultVisibility(t *testing.T) { + ctx := t.Context() + repos, auth, _, _ := newAuthTestServices(t) + user, _, err := auth.Register(ctx, "viewer", "secret-pass") + if err != nil { + t.Fatalf("register user: %v", err) + } + if err := repos.DB.Create(&model.TelegramBinding{ + TelegramUserID: 30001, + TelegramName: "@viewer", + ChatID: 30001, + UserID: user.ID, + }).Error; err != nil { + t.Fatalf("create binding: %v", err) + } + if err := repos.DB.AutoMigrate(&model.NotifyChannel{}); err != nil { + t.Fatalf("migrate notify_channels: %v", err) + } + // 配置一个绑定该 Telegram 用户的渠道(无 bot_token,避免测试触发网络请求)。 + cfg, _ := json.Marshal(map[string]string{"admin_user_ids": "30001"}) + if err := repos.DB.Create(&model.NotifyChannel{ + Name: "Telegram", + Type: "telegram", + Enabled: true, + Config: string(cfg), + }).Error; err != nil { + t.Fatalf("create channel: %v", err) + } + + before, err := repos.User.FindByID(ctx, user.ID) + if err != nil || before == nil { + t.Fatalf("load user before toggle: %v", err) + } + + bot := NewTelegramBotService(zap.NewNop(), repos, nil) + update, _ := json.Marshal(TelegramUpdate{ + UpdateID: 1, + CallbackQuery: &TelegramCallbackQuery{ + ID: "cb1", + From: TelegramUser{ID: 30001, Username: "viewer", FirstName: "Viewer"}, + Message: &TelegramMessage{MessageID: 5, Chat: TelegramChat{ID: 30001, Type: "private"}}, + Data: "adult_toggle", + }, + }) + // reply 因 bot_token 为空会返回错误,但成人目录状态应已在数据库中被切换。 + _ = bot.HandleWebhook(ctx, update) + + updated, err := repos.User.FindByID(ctx, user.ID) + if err != nil || updated == nil { + t.Fatalf("reload user: %v", err) + } + if updated.HideAdult == before.HideAdult { + t.Fatalf("adult_toggle callback should have flipped HideAdult (was %v)", before.HideAdult) + } +} + func TestTelegramStartClearsStaleUserBinding(t *testing.T) { repos, _, _, _ := newAuthTestServices(t) if err := repos.DB.Create(&model.TelegramBinding{ diff --git a/internal/service/transfer.go b/internal/service/transfer.go new file mode 100644 index 0000000..e492bd4 --- /dev/null +++ b/internal/service/transfer.go @@ -0,0 +1,98 @@ +// Package service — file transfer strategies for the organizer. +// +// 整理媒体时支持四种转移方式: +// +// move 移动(同盘 rename,跨盘 copy+删除源)——会移除源文件 +// copy 复制(保留源文件) +// hardlink 硬链接(同盘零额外占用,保留源文件;做种不受影响) +// symlink 软链接(保留源文件,指向源) +// +// 除 move 外,其余方式都保留源文件,因此 qBittorrent 等下载器仍能在原 +// 路径找到数据继续做种上传。 +package service + +import ( + "fmt" + "io" + "os" + "path/filepath" + "strings" +) + +// TransferMode 表示整理时文件的转移方式。 +type TransferMode string + +const ( + // TransferMove 移动:同盘 rename,跨盘 copy+删除源。 + TransferMove TransferMode = "move" + // TransferCopy 复制:保留源文件。 + TransferCopy TransferMode = "copy" + // TransferHardlink 硬链接:同盘零额外占用并保留源文件,做种不受影响。 + TransferHardlink TransferMode = "hardlink" + // TransferSymlink 软链接:保留源文件,目标指向源。 + TransferSymlink TransferMode = "symlink" +) + +// parseTransferMode 解析转移方式字符串,无法识别时回退为默认的移动。 +func parseTransferMode(s string) TransferMode { + switch strings.ToLower(strings.TrimSpace(s)) { + case "copy", "复制": + return TransferCopy + case "hardlink", "hard", "link", "硬链接", "硬连接": + return TransferHardlink + case "symlink", "soft", "softlink", "软链接", "软连接", "符号链接": + return TransferSymlink + default: + return TransferMove + } +} + +// keepsSource 报告该转移方式是否会保留源文件(用于做种)。 +func (m TransferMode) keepsSource() bool { + return m == TransferCopy || m == TransferHardlink || m == TransferSymlink +} + +// transferFile 按指定方式把 src 转移到 dst。 +// dst 已存在时一律报错,绝不覆盖(防止不同 release 改名后互相覆盖)。 +func transferFile(src, dst string, mode TransferMode) error { + if _, err := os.Stat(dst); err == nil { + return fmt.Errorf("destination already exists: %s", dst) + } + switch mode { + case TransferCopy: + return copyFile(src, dst) + case TransferHardlink: + if err := os.Link(src, dst); err != nil { + // 跨文件系统无法硬链接,退化为复制,仍保留源文件以便继续做种。 + return copyFile(src, dst) + } + return nil + case TransferSymlink: + target := src + if abs, err := filepath.Abs(src); err == nil { + target = abs + } + return os.Symlink(target, dst) + default: // TransferMove + return moveFile(src, dst) + } +} + +// copyFile 流式复制 src 到 dst(保留源文件)。O_EXCL 保证不覆盖已存在目标。 +func copyFile(src, dst string) error { + in, err := os.Open(src) + if err != nil { + return err + } + defer in.Close() + f, err := os.OpenFile(dst, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0o644) + if err != nil { + return err + } + if _, werr := io.Copy(f, in); werr != nil { + f.Close() + os.Remove(dst) + return werr + } + return f.Close() +} diff --git a/internal/service/transfer_test.go b/internal/service/transfer_test.go new file mode 100644 index 0000000..8d1021e --- /dev/null +++ b/internal/service/transfer_test.go @@ -0,0 +1,123 @@ +package service + +import ( + "os" + "path/filepath" + "testing" +) + +func TestParseTransferMode(t *testing.T) { + cases := map[string]TransferMode{ + "": TransferMove, + "move": TransferMove, + "移动": TransferMove, + "copy": TransferCopy, + "复制": TransferCopy, + "hardlink": TransferHardlink, + "硬链接": TransferHardlink, + "symlink": TransferSymlink, + "软链接": TransferSymlink, + "garbage": TransferMove, + } + for in, want := range cases { + if got := parseTransferMode(in); got != want { + t.Errorf("parseTransferMode(%q) = %q, want %q", in, got, want) + } + } +} + +func writeTemp(t *testing.T, dir, name, content string) string { + t.Helper() + p := filepath.Join(dir, name) + if err := os.WriteFile(p, []byte(content), 0o644); err != nil { + t.Fatalf("write %s: %v", p, err) + } + return p +} + +func TestTransferFileCopyKeepsSource(t *testing.T) { + dir := t.TempDir() + src := writeTemp(t, dir, "src.mkv", "payload") + dst := filepath.Join(dir, "out", "dst.mkv") + if err := os.MkdirAll(filepath.Dir(dst), 0o755); err != nil { + t.Fatal(err) + } + if err := transferFile(src, dst, TransferCopy); err != nil { + t.Fatalf("copy: %v", err) + } + if _, err := os.Stat(src); err != nil { + t.Fatalf("copy must keep source: %v", err) + } + b, _ := os.ReadFile(dst) + if string(b) != "payload" { + t.Fatalf("copied content = %q", b) + } +} + +func TestTransferFileHardlinkSharesInodeAndKeepsSource(t *testing.T) { + dir := t.TempDir() + src := writeTemp(t, dir, "src.mkv", "payload") + dst := filepath.Join(dir, "dst.mkv") + if err := transferFile(src, dst, TransferHardlink); err != nil { + t.Fatalf("hardlink: %v", err) + } + si, err := os.Stat(src) + if err != nil { + t.Fatalf("hardlink must keep source: %v", err) + } + di, err := os.Stat(dst) + if err != nil { + t.Fatalf("stat dst: %v", err) + } + if !os.SameFile(si, di) { + t.Fatal("hardlink dst should share inode with source") + } +} + +func TestTransferFileSymlinkKeepsSource(t *testing.T) { + dir := t.TempDir() + src := writeTemp(t, dir, "src.mkv", "payload") + dst := filepath.Join(dir, "dst.mkv") + if err := transferFile(src, dst, TransferSymlink); err != nil { + t.Fatalf("symlink: %v", err) + } + fi, err := os.Lstat(dst) + if err != nil { + t.Fatalf("lstat dst: %v", err) + } + if fi.Mode()&os.ModeSymlink == 0 { + t.Fatal("dst should be a symlink") + } + if _, err := os.Stat(src); err != nil { + t.Fatalf("symlink must keep source: %v", err) + } +} + +func TestTransferFileMoveRemovesSource(t *testing.T) { + dir := t.TempDir() + src := writeTemp(t, dir, "src.mkv", "payload") + dst := filepath.Join(dir, "dst.mkv") + if err := transferFile(src, dst, TransferMove); err != nil { + t.Fatalf("move: %v", err) + } + if _, err := os.Stat(src); !os.IsNotExist(err) { + t.Fatalf("move should remove source, stat err = %v", err) + } + if _, err := os.Stat(dst); err != nil { + t.Fatalf("move should create dst: %v", err) + } +} + +func TestTransferFileNeverOverwrites(t *testing.T) { + dir := t.TempDir() + src := writeTemp(t, dir, "src.mkv", "new") + dst := writeTemp(t, dir, "dst.mkv", "existing") + for _, mode := range []TransferMode{TransferMove, TransferCopy, TransferHardlink, TransferSymlink} { + if err := transferFile(src, dst, mode); err == nil { + t.Fatalf("mode %q should refuse to overwrite existing dst", mode) + } + if b, _ := os.ReadFile(dst); string(b) != "existing" { + t.Fatalf("mode %q clobbered existing dst", mode) + } + } +} diff --git a/internal/service/watcher.go b/internal/service/watcher.go index c24a3aa..d7f5f1d 100644 --- a/internal/service/watcher.go +++ b/internal/service/watcher.go @@ -1,15 +1,20 @@ // Package service — filesystem watcher. // // WatcherService observes every enabled library root with fsnotify and -// debounces incoming events into per-library re-scans. New / renamed +// debounces incoming events into incremental, per-file ingests. New / renamed // files become Media rows; deletes remove them. // +// 设计目标:只在「有新增/变更媒体」时增量入库,绝不因为单个文件变化就对整个 +// 媒体库做全量重扫——全量重扫会反复读盘、损伤硬盘,也是用户明确要避免的。 +// 因此 watcher 递归监听库内所有子目录,事件去抖后只处理具体变化的路径。 +// // The watcher runs in the background and is started after migrations // complete. It survives library add / delete via Refresh(). package service import ( "context" + "os" "path/filepath" "sync" "time" @@ -20,6 +25,13 @@ import ( "github.com/ShukeBta/MediaStationGo/internal/repository" ) +// pendingEvent records the most recent change to a path and the library it +// belongs to, for debounced incremental processing. +type pendingEvent struct { + libraryID string + ts time.Time +} + // WatcherService is a thin orchestrator on top of fsnotify. type WatcherService struct { log *zap.Logger @@ -28,8 +40,8 @@ type WatcherService struct { mu sync.Mutex watcher *fsnotify.Watcher - watched map[string]string // dir -> libraryID - pending map[string]time.Time + watched map[string]string // dir -> libraryID + pending map[string]pendingEvent // path -> most recent change stop chan struct{} } @@ -40,7 +52,7 @@ func NewWatcherService(log *zap.Logger, repo *repository.Container, scanner *Sca repo: repo, scanner: scanner, watched: make(map[string]string), - pending: make(map[string]time.Time), + pending: make(map[string]pendingEvent), stop: make(chan struct{}), } } @@ -79,12 +91,17 @@ func (w *WatcherService) Refresh(ctx context.Context) error { w.mu.Lock() defer w.mu.Unlock() + // Map every directory (root + all subdirectories) to its library so new + // files anywhere in the tree raise events — fsnotify itself is + // non-recursive, so we register each directory explicitly. current := make(map[string]string) for _, l := range libs { if !l.Enabled { continue } - current[l.Path] = l.ID + for _, dir := range listDirsForWatch(l.Path) { + current[dir] = l.ID + } } // Remove disappeared paths. for path := range w.watched { @@ -93,7 +110,7 @@ func (w *WatcherService) Refresh(ctx context.Context) error { delete(w.watched, path) } } - // Add new ones (top-level only — fsnotify is non-recursive). + // Add new ones. for path, id := range current { if _, ok := w.watched[path]; ok { continue @@ -107,6 +124,34 @@ func (w *WatcherService) Refresh(ctx context.Context) error { return nil } +// listDirsForWatch returns root plus every (non-hidden) subdirectory so the +// watcher can register the whole tree recursively. +func listDirsForWatch(root string) []string { + dirs := []string{root} + _ = walk(root, func(path string, info walkInfo) error { + if info.isDir && path != root { + dirs = append(dirs, path) + } + return nil + }) + return dirs +} + +// watchDirRecursive registers a newly-created directory subtree so files +// copied into it afterwards still raise events. +func (w *WatcherService) watchDirRecursive(dir, libraryID string) { + for _, d := range listDirsForWatch(dir) { + if _, ok := w.watched[d]; ok { + continue + } + if err := w.watcher.Add(d); err != nil { + w.log.Debug("watch add (recursive) failed", zap.String("path", d), zap.Error(err)) + continue + } + w.watched[d] = libraryID + } +} + // loop drains fsnotify events and pushes the affected library into the // pending map. The actual rescan happens in the debouncer goroutine. func (w *WatcherService) loop(ctx context.Context) { @@ -130,8 +175,16 @@ func (w *WatcherService) loop(ctx context.Context) { if lib == "" { continue } + // 新建目录:立即递归纳入监听,确保随后拷入的文件也能触发事件。 + if ev.Op&fsnotify.Create != 0 { + if fi, err := os.Stat(ev.Name); err == nil && fi.IsDir() { + w.mu.Lock() + w.watchDirRecursive(ev.Name, lib) + w.mu.Unlock() + } + } w.mu.Lock() - w.pending[lib] = time.Now() + w.pending[ev.Name] = pendingEvent{libraryID: lib, ts: time.Now()} w.mu.Unlock() case err, ok := <-w.watcher.Errors: if !ok { @@ -160,9 +213,17 @@ func (w *WatcherService) findLibrary(path string) string { } } -// debouncer drains the pending set every 5 s and triggers a rescan per -// library. Coalescing avoids storming the disk on bulk operations -// (mass-rename, large copies). +// duePath couples a settled path with its library for incremental processing. +type duePath struct { + path string + libraryID string +} + +// debouncer drains the pending set every 5 s and processes each settled path +// incrementally: existing files are ingested (single-file upsert), vanished +// files are removed. Coalescing by path avoids storming the disk during bulk +// operations (mass-rename, large copies), and crucially we never re-walk the +// entire library — only the paths that actually changed. func (w *WatcherService) debouncer(ctx context.Context) { t := time.NewTicker(5 * time.Second) defer t.Stop() @@ -175,20 +236,39 @@ func (w *WatcherService) debouncer(ctx context.Context) { case <-t.C: } w.mu.Lock() - due := make([]string, 0, len(w.pending)) + due := make([]duePath, 0, len(w.pending)) now := time.Now() - for id, ts := range w.pending { - if now.Sub(ts) >= 5*time.Second { - due = append(due, id) - delete(w.pending, id) + for path, ev := range w.pending { + if now.Sub(ev.ts) >= 5*time.Second { + due = append(due, duePath{path: path, libraryID: ev.libraryID}) + delete(w.pending, path) } } w.mu.Unlock() - for _, id := range due { - w.log.Info("watcher triggered rescan", zap.String("library_id", id)) - if _, err := w.scanner.ScanLibrary(ctx, id); err != nil { - w.log.Warn("watcher rescan failed", zap.Error(err)) - } + for _, d := range due { + w.process(ctx, d) } } } + +// process ingests or removes a single changed path. +func (w *WatcherService) process(ctx context.Context, d duePath) { + fi, err := os.Stat(d.path) + if err != nil { + // Vanished (delete/rename away): drop its media row if any. + if removed, derr := w.scanner.RemovePath(ctx, d.path); derr != nil { + w.log.Warn("watcher remove failed", zap.String("path", d.path), zap.Error(derr)) + } else if removed > 0 { + w.log.Info("watcher removed media", zap.String("path", d.path)) + } + return + } + if fi.IsDir() { + return // directory events only matter for registering new watches + } + if added, ierr := w.scanner.IngestPath(ctx, d.libraryID, d.path); ierr != nil { + w.log.Warn("watcher ingest failed", zap.String("path", d.path), zap.Error(ierr)) + } else if added { + w.log.Info("watcher ingested media", zap.String("path", d.path)) + } +} diff --git a/web/src/api/tools.ts b/web/src/api/tools.ts index 607d50d..0817bba 100644 --- a/web/src/api/tools.ts +++ b/web/src/api/tools.ts @@ -3,15 +3,25 @@ import { api } from './client' // toolsAPI groups admin-only endpoints that don't fit the other domain // modules: organizing media files into the canonical naming layout, and // dispatching a test notification through the configured channels. +// OrganizeOverrides are optional single-request overrides for an organize +// action. Empty fields fall back to the system settings. +export interface OrganizeOverrides { + target_path?: string + transfer_mode?: string +} + export const toolsAPI = { - organizeMedia: (mediaID: string) => + organizeMedia: (mediaID: string, opts?: OrganizeOverrides) => api - .post<{ path: string }>(`/admin/media/${mediaID}/organize`) + .post<{ path: string }>(`/admin/media/${mediaID}/organize`, opts ?? {}) .then((r) => r.data), - organizeLibrary: (libraryID: string) => + organizeLibrary: (libraryID: string, opts?: OrganizeOverrides) => api - .post>(`/admin/libraries/${libraryID}/organize`) + .post>( + `/admin/libraries/${libraryID}/organize`, + opts ?? {}, + ) .then((r) => r.data), notifyTest: (title: string, body: string) => diff --git a/web/src/pages/SettingsPage.tsx b/web/src/pages/SettingsPage.tsx index b8a69eb..d844b77 100644 --- a/web/src/pages/SettingsPage.tsx +++ b/web/src/pages/SettingsPage.tsx @@ -158,6 +158,24 @@ const GROUPS: SettingGroup[] = [ hint: '留空则默认整理到各媒体库对应路径(见下方参考)', placeholder: '/mnt/media/organized', }, + { + key: 'organize.transfer_mode', + label: '默认转移方式', + type: 'select', + hint: '移动会删除源文件;复制/硬链接/软链接保留源文件,PT 做种不中断。硬链接同盘零额外占用。', + options: [ + { value: 'move', label: '移动(删除源文件)' }, + { value: 'copy', label: '复制(保留源文件)' }, + { value: 'hardlink', label: '硬链接(保留源,做种不中断,不占双倍空间)' }, + { value: 'symlink', label: '软链接(符号链接,保留源)' }, + ], + }, + { + key: 'organize.keep_seeding', + label: '保种(整理后继续做种上传)', + type: 'toggle', + hint: '开启后即使选择「移动」也会自动改用硬链接(跨盘退化为复制)保留源文件,确保 qBittorrent 转移后继续做种上传。', + }, { key: 'organize.movie_format', label: '电影命名格式', @@ -195,6 +213,12 @@ const GROUPS: SettingGroup[] = [ type: 'text', placeholder: 'zh-CN', }, + { + key: 'scan.periodic_enabled', + label: '周期性整库重扫', + type: 'toggle', + hint: '默认关闭。文件新增/变更由实时监听增量入库,无需定时全量重扫。开启会每 60 分钟重扫整库,频繁读盘会损伤硬盘,一般无需开启。', + }, ], }, { diff --git a/web/src/pages/ToolsPage.tsx b/web/src/pages/ToolsPage.tsx index e5fbaab..a29e32b 100644 --- a/web/src/pages/ToolsPage.tsx +++ b/web/src/pages/ToolsPage.tsx @@ -61,6 +61,17 @@ function OrganizePanel() { const [smartClassify, setSmartClassify] = useState(false) const [loadingSettings, setLoadingSettings] = useState(true) + // 单次整理覆盖项:留空则沿用设置页的默认整理目录与转移方式。 + const [targetPath, setTargetPath] = useState('') + const [transferMode, setTransferMode] = useState('') + + const overrides = () => { + const o: { target_path?: string; transfer_mode?: string } = {} + if (targetPath.trim()) o.target_path = targetPath.trim() + if (transferMode) o.transfer_mode = transferMode + return o + } + useEffect(() => { libraryAPI.list().then(setLibraries).catch(() => undefined) }, []) @@ -87,7 +98,7 @@ function OrganizePanel() { if (!libraryID) return setRunning(true) try { - await toolsAPI.organizeLibrary(libraryID) + await toolsAPI.organizeLibrary(libraryID, overrides()) toast.success('已触发媒体库整理') } catch (err: unknown) { const msg = @@ -116,8 +127,8 @@ function OrganizePanel() { const onOrganizeOne = async (m: Media) => { setBusyID(m.id) try { - const r = await toolsAPI.organizeMedia(m.id) - toast.success(`已移动到 ${r.path}`) + const r = await toolsAPI.organizeMedia(m.id, overrides()) + toast.success(`已整理到 ${r.path}`) } catch (err: unknown) { const msg = (err as { response?: { data?: { error?: string } } })?.response?.data?.error ?? @@ -151,6 +162,32 @@ function OrganizePanel() { )} +
+ + +
+