mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-07 13:56:37 +08:00
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>
This commit is contained in:
+100
-20
@@ -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))
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user