Files
MeBox/internal/service/watcher.go
T
Kiro 7bd6bcfb1b feat: downloads, RSS, subtitles, Bangumi, watcher, stats, profile, audit
Backend
  - service/bangumi.go: Bangumi (bgm.tv) scraper for anime libraries.
  - service/episode_parser.go: SxxExx / NxE / EPxx / 第NN集 parser with
    unit tests; consumed by the scanner for tv/anime libraries.
  - service/scraper.go: provider chain orchestrator picks Bangumi for
    anime libraries (TMDb fallback), TMDb for everything else; emits
    'no_match' rows so we don't retry forever; AnyEnabled() for the
    scanner kick.
  - service/scanner.go: writes season/episode numbers for tv/anime libs
    and reports a 'probed' counter alongside 'added'.
  - service/transcoder.go (existing): unchanged, stays per-media.
  - service/subtitle.go: discovers external subtitles next to the source
    file (.srt/.vtt/.ass/.ssa) and converts SRT/ASS to WebVTT on the fly.
  - service/qbittorrent.go: thread-safe qBittorrent v2 Web UI client
    (login / add / list / delete) with cookie-jar reuse.
  - service/downloads.go: persists download tasks, reads runtime
    qbittorrent.* settings, polls /torrents/info every 5 s and pushes
    the result to WS subscribers ('download' topic).
  - service/subscription.go: 10-minute RSS poller with regex filter,
    GUID dedup persisted in the settings table, and 'subscription' WS
    events on enqueue.
  - service/watcher.go: fsnotify watcher with 5 s coalescing debouncer
    that triggers per-library rescans on create/rename/remove.
  - service/stats.go: dashboard snapshot — totals, recently added,
    gopsutil-driven CPU/mem/disk readings.
  - service/profile.go: non-credential profile patch + admin role mutator.
  - service/audit.go: best-effort writer for the access_logs table.
  - service/service.go: container wires every new service; Boot() spins
    up watcher / downloads poller / subscription scheduler; Close()
    tears them down on graceful shutdown.

Handlers
  - new files: downloads.go, subscriptions.go, subtitles.go, series.go,
    stats.go, profile.go, util.go.
  - handler.go: registers PATCH /me, /libraries/:id/seasons,
    /media/:id/subtitles, /subtitles/:id, /downloads*, /subscriptions*,
    /stats, /admin/users/:id/role.
  - auth.go / media.go: write audit rows for login + library CRUD and
    refresh the watcher when libraries change.

Frontend
  - api: new helpers for downloads, subscriptions, profile, series, stats,
    subtitles; library helper gained scrape().
  - hooks/useWebSocket.ts: shared connection with 3 s reconnect.
  - components/GlobalEvents.tsx: surfaces scan / scrape / subscription
    completion as toasts (mounted at app root).
  - pages: Library now switches to a season-grouped layout for tv/anime
    libraries; Player attaches WebVTT <track> elements; new pages for
    Downloads (live torrent table), Subscriptions, Profile, Stats.
  - components/Layout.tsx + App.tsx: sidebar groups (媒体库 / 自动化 /
    账号 / 管理) and routes for the new pages; /stats and /admin remain
    admin-only.
  - types/index.ts: new types — Subscription, DownloadTask, QBitTorrent,
    Hardware, StatsSnapshot.

Verified: go build, go vet, go test (incl. ParseEpisode + srtToVTT +
stripASSTags) all pass; tsc -b && vite build emits 17 route chunks plus
the deferred hls chunk; main bundle 247 KB / 83 KB gzipped.
2026-05-14 16:07:49 +00:00

195 lines
4.5 KiB
Go

// Package service — filesystem watcher.
//
// WatcherService observes every enabled library root with fsnotify and
// debounces incoming events into per-library re-scans. New / renamed
// files become Media rows; deletes remove them.
//
// The watcher runs in the background and is started after migrations
// complete. It survives library add / delete via Refresh().
package service
import (
"context"
"path/filepath"
"sync"
"time"
"github.com/fsnotify/fsnotify"
"go.uber.org/zap"
"github.com/ShukeBta/MediaStationGo/internal/repository"
)
// WatcherService is a thin orchestrator on top of fsnotify.
type WatcherService struct {
log *zap.Logger
repo *repository.Container
scanner *ScannerService
mu sync.Mutex
watcher *fsnotify.Watcher
watched map[string]string // dir -> libraryID
pending map[string]time.Time
stop chan struct{}
}
// NewWatcherService is the constructor.
func NewWatcherService(log *zap.Logger, repo *repository.Container, scanner *ScannerService) *WatcherService {
return &WatcherService{
log: log,
repo: repo,
scanner: scanner,
watched: make(map[string]string),
pending: make(map[string]time.Time),
stop: make(chan struct{}),
}
}
// Start initialises the underlying fsnotify watcher and registers every
// library root currently in the database.
func (w *WatcherService) Start(ctx context.Context) error {
fw, err := fsnotify.NewWatcher()
if err != nil {
return err
}
w.watcher = fw
if err := w.Refresh(ctx); err != nil {
w.log.Warn("watcher refresh failed", zap.Error(err))
}
go w.loop(ctx)
go w.debouncer(ctx)
return nil
}
// Stop tears down the watcher (called on graceful shutdown).
func (w *WatcherService) Stop() {
close(w.stop)
if w.watcher != nil {
_ = w.watcher.Close()
}
}
// Refresh reads the library list and adjusts the set of watched
// directories. Idempotent — safe to call after every CRUD.
func (w *WatcherService) Refresh(ctx context.Context) error {
libs, err := w.repo.Library.List(ctx)
if err != nil {
return err
}
w.mu.Lock()
defer w.mu.Unlock()
current := make(map[string]string)
for _, l := range libs {
if !l.Enabled {
continue
}
current[l.Path] = l.ID
}
// Remove disappeared paths.
for path := range w.watched {
if _, ok := current[path]; !ok {
_ = w.watcher.Remove(path)
delete(w.watched, path)
}
}
// Add new ones (top-level only — fsnotify is non-recursive).
for path, id := range current {
if _, ok := w.watched[path]; ok {
continue
}
if err := w.watcher.Add(path); err != nil {
w.log.Warn("watch add failed", zap.String("path", path), zap.Error(err))
continue
}
w.watched[path] = id
}
return nil
}
// 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) {
if w.watcher == nil {
return
}
for {
select {
case <-ctx.Done():
return
case <-w.stop:
return
case ev, ok := <-w.watcher.Events:
if !ok {
return
}
if ev.Op&(fsnotify.Create|fsnotify.Remove|fsnotify.Rename|fsnotify.Write) == 0 {
continue
}
lib := w.findLibrary(ev.Name)
if lib == "" {
continue
}
w.mu.Lock()
w.pending[lib] = time.Now()
w.mu.Unlock()
case err, ok := <-w.watcher.Errors:
if !ok {
return
}
w.log.Warn("watcher error", zap.Error(err))
}
}
}
// findLibrary maps a path back to the watching library ID, taking the
// shortest matching prefix.
func (w *WatcherService) findLibrary(path string) string {
w.mu.Lock()
defer w.mu.Unlock()
dir := filepath.Dir(path)
for {
if id, ok := w.watched[dir]; ok {
return id
}
parent := filepath.Dir(dir)
if parent == dir {
return ""
}
dir = parent
}
}
// 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).
func (w *WatcherService) debouncer(ctx context.Context) {
t := time.NewTicker(5 * time.Second)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-w.stop:
return
case <-t.C:
}
w.mu.Lock()
due := make([]string, 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)
}
}
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))
}
}
}
}