mirror of
https://github.com/truewhile/MeBox.git
synced 2026-09-30 03:36:37 +08:00
7bd6bcfb1b
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.
195 lines
4.5 KiB
Go
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))
|
|
}
|
|
}
|
|
}
|
|
}
|