stremio/node/media-node/internal/scanner/scanner.go
Jos Vooges | STH 18eb702de2 Add Synology Download Station search and 4K import pipeline.
Admin can search DS, confirm TMDB/IMDb, then auto-move finished downloads into a fixed 4K library folder via the media-node.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-30 17:42:12 +02:00

745 lines
18 KiB
Go

package scanner
import (
"log"
"os"
"path/filepath"
"strings"
"sync"
"time"
"github.com/fsnotify/fsnotify"
"github.com/sthmedia/media-node/internal/config"
"github.com/sthmedia/media-node/internal/database"
)
var videoExtensions = map[string]bool{
".mkv": true, ".mp4": true, ".m4v": true,
".avi": true, ".ts": true, ".m2ts": true,
}
type MediaFileInfo struct {
LocalFileID string `json:"localFileId"`
SizeBytes int64 `json:"sizeBytes"`
Container string `json:"container,omitempty"`
Resolution string `json:"resolution,omitempty"`
VideoCodec string `json:"videoCodec,omitempty"`
AudioCodec string `json:"audioCodec,omitempty"`
ReleaseName string `json:"releaseName"`
ModifiedAt string `json:"modifiedAt"`
MediaType string `json:"mediaType"`
ShelfID string `json:"shelfId,omitempty"`
ParsedMovie map[string]interface{} `json:"parsedMovie,omitempty"`
ParsedEpisode map[string]interface{} `json:"parsedEpisode,omitempty"`
}
type LibraryEvent struct {
Type string `json:"type"`
NodeRevision int64 `json:"nodeRevision"`
File MediaFileInfo `json:"file"`
}
type EventHandler func(events []LibraryEvent)
type Scanner struct {
cfg *config.Config
store *database.Store
onEvents EventHandler
onScanComplete func()
mu sync.Mutex
scanning bool
rescanQueued bool
status string
resetSchedule chan struct{}
reloadWatch chan struct{}
}
func New(cfg *config.Config, store *database.Store, onEvents EventHandler) *Scanner {
return &Scanner{
cfg: cfg,
store: store,
onEvents: onEvents,
status: "idle",
resetSchedule: make(chan struct{}, 2),
reloadWatch: make(chan struct{}, 1),
}
}
// SetOnScanComplete registers a hook after a full scan finishes (no pending re-run).
func (s *Scanner) SetOnScanComplete(fn func()) {
s.mu.Lock()
s.onScanComplete = fn
s.mu.Unlock()
}
func (s *Scanner) Status() string {
s.mu.Lock()
defer s.mu.Unlock()
return s.status
}
// WaitUntilIdle blocks until no full scan is running or queued.
func (s *Scanner) WaitUntilIdle() {
for {
s.mu.Lock()
busy := s.scanning || s.rescanQueued
s.mu.Unlock()
if !busy {
return
}
time.Sleep(500 * time.Millisecond)
}
}
func (s *Scanner) Start() {
// One inventory pass at boot, then realtime watch + scheduled catch-up.
go s.runStartupScan()
go s.runNightlyFullScan()
go s.runIntervalFullScan()
go s.runIncrementalCatchUp()
go s.watchFilesystem()
}
// UpdateSchedule applies admin/master scan settings and restarts timers.
func (s *Scanner) UpdateSchedule(interval, at string) {
s.mu.Lock()
if interval != "" {
s.cfg.Scanner.FullScanInterval = interval
}
if at != "" {
s.cfg.Scanner.FullScanAt = at
}
s.mu.Unlock()
select {
case s.resetSchedule <- struct{}{}:
default:
}
select {
case s.resetSchedule <- struct{}{}:
default:
}
log.Printf("Scanner schedule updated: interval=%s at=%s", s.cfg.Scanner.FullScanInterval, s.cfg.Scanner.FullScanAt)
}
// UpdateMediaPaths replaces movie/series roots, rebuilds fsnotify watchers, and rescans.
func (s *Scanner) UpdateMediaPaths(movies, series []string) {
roots := make([]config.MediaRoot, 0, len(movies)+len(series))
for _, p := range movies {
roots = append(roots, config.MediaRoot{Path: p, Kind: "movie"})
}
for _, p := range series {
roots = append(roots, config.MediaRoot{Path: p, Kind: "episode"})
}
s.UpdateMediaRoots(roots)
}
// UpdateMediaRoots replaces scan roots (with optional shelf IDs).
// Returns true when roots changed (and a full scan was started).
func (s *Scanner) UpdateMediaRoots(roots []config.MediaRoot) bool {
s.mu.Lock()
if mediaRootsEqual(s.cfg.Media.EffectiveRoots(), roots) {
s.mu.Unlock()
log.Printf("Media roots unchanged (%d roots) — skip rescan", len(roots))
return false
}
s.cfg.Media.Roots = append([]config.MediaRoot{}, roots...)
var movies, series []string
for _, r := range roots {
if r.Kind == "episode" || r.Kind == "series" {
series = append(series, r.Path)
} else {
movies = append(movies, r.Path)
}
}
s.cfg.Media.Movies = movies
s.cfg.Media.Series = series
s.mu.Unlock()
select {
case s.reloadWatch <- struct{}{}:
default:
}
log.Printf("Media roots updated: %d roots — starting full scan", len(roots))
go s.FullScan()
return true
}
func mediaRootsEqual(a, b []config.MediaRoot) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if filepath.Clean(a[i].Path) != filepath.Clean(b[i].Path) ||
a[i].ShelfID != b[i].ShelfID ||
a[i].Kind != b[i].Kind {
return false
}
}
return true
}
func (s *Scanner) runStartupScan() {
log.Println("Startup library inventory scan...")
s.FullScan()
}
func (s *Scanner) runNightlyFullScan() {
for {
s.mu.Lock()
hour, minute, ok := s.cfg.FullScanClock()
s.mu.Unlock()
if !ok {
select {
case <-s.resetSchedule:
continue
}
}
next := nextLocalClock(time.Now(), hour, minute)
log.Printf("Next nightly full scan at %s", next.Format(time.RFC3339))
timer := time.NewTimer(time.Until(next))
select {
case <-timer.C:
log.Println("Nightly catch-up full scan starting...")
s.FullScan()
case <-s.resetSchedule:
if !timer.Stop() {
select {
case <-timer.C:
default:
}
}
}
}
}
func (s *Scanner) runIntervalFullScan() {
for {
s.mu.Lock()
d := s.cfg.FullScanDuration()
s.mu.Unlock()
if d <= 0 {
select {
case <-s.resetSchedule:
continue
}
}
log.Printf("Interval full scan every %s", d)
timer := time.NewTimer(d)
select {
case <-timer.C:
s.FullScan()
case <-s.resetSchedule:
if !timer.Stop() {
select {
case <-timer.C:
default:
}
}
}
}
}
func nextLocalClock(now time.Time, hour, minute int) time.Time {
loc := now.Location()
next := time.Date(now.Year(), now.Month(), now.Day(), hour, minute, 0, 0, loc)
if !next.After(now) {
next = next.Add(24 * time.Hour)
}
return next
}
const incrementalCatchUpInterval = 10 * time.Minute
// runIncrementalCatchUp picks up new/changed files when fsnotify misses events
// (common with SMB uploads or inotify limits on large libraries).
func (s *Scanner) runIncrementalCatchUp() {
timer := time.NewTimer(incrementalCatchUpInterval)
defer timer.Stop()
for {
<-timer.C
s.incrementalScan()
timer.Reset(incrementalCatchUpInterval)
}
}
func (s *Scanner) incrementalScan() {
s.mu.Lock()
if s.scanning {
s.mu.Unlock()
return
}
roots := s.cfg.Media.EffectiveRoots()
s.mu.Unlock()
if len(roots) == 0 {
return
}
newCount := 0
for _, root := range roots {
kind := root.Kind
if kind == "series" {
kind = "episode"
}
if root.Path == "" {
continue
}
_ = filepath.Walk(root.Path, func(path string, info os.FileInfo, err error) error {
if err != nil || info.IsDir() || !isVideo(path) {
return nil
}
existing, _ := s.store.GetFileByPath(path)
if existing != nil &&
existing.ModifiedAt.Equal(info.ModTime()) &&
existing.SizeBytes == info.Size() {
return nil
}
eventType := "FILE_ADDED"
if existing != nil {
eventType = "FILE_UPDATED"
}
shelf := shelfForPath(path, roots)
s.indexFile(path, info, eventType, kind, shelf, true)
newCount++
return nil
})
}
if newCount > 0 {
log.Printf("Incremental scan indexed %d new/changed file(s)", newCount)
}
}
func (s *Scanner) FullScan() {
for {
s.mu.Lock()
if s.scanning {
s.rescanQueued = true
s.mu.Unlock()
log.Println("Full scan already running — queued re-run after current pass")
return
}
s.scanning = true
s.rescanQueued = false
s.status = "scanning"
roots := s.cfg.Media.EffectiveRoots()
s.mu.Unlock()
log.Printf("Starting full media scan (%d roots)...", len(roots))
found := make(map[string]bool)
aborted := false
for _, root := range roots {
s.mu.Lock()
if s.rescanQueued {
s.mu.Unlock()
aborted = true
log.Printf("Aborting scan early (roots/config changed); will re-run")
break
}
s.mu.Unlock()
kind := root.Kind
if kind == "series" {
kind = "episode"
}
before := len(found)
log.Printf("Scanning root %s (%s) shelf=%s ...", root.Path, kind, root.ShelfID)
s.walkRoot(root.Path, kind, root.ShelfID, found, false)
log.Printf("Scan root done %s → +%d files (total %d)", root.Path, len(found)-before, len(found))
}
if aborted {
s.mu.Lock()
s.scanning = false
s.status = "idle"
s.mu.Unlock()
log.Println("Re-running full scan after path/config change...")
continue
}
// Local prune only — Master learns removals via the full library sync that follows.
existing, _ := s.store.AllFiles()
removed := 0
for _, f := range existing {
if !found[f.Path] {
_ = s.store.DeleteFile(f.LocalFileID)
removed++
}
}
log.Printf("Full scan completed (%d files, %d removed locally)", len(found), removed)
s.mu.Lock()
if s.rescanQueued {
s.scanning = false
s.status = "idle"
s.mu.Unlock()
log.Println("Re-running full scan after path/config change...")
continue
}
s.scanning = false
s.status = "idle"
doneHook := s.onScanComplete
s.mu.Unlock()
if doneHook != nil {
go doneHook()
}
return
}
}
func (s *Scanner) walkRoot(dir, mediaType, shelfID string, found map[string]bool, emitEvents bool) {
if dir == "" {
return
}
filepath.Walk(dir, func(path string, info os.FileInfo, err error) error {
if err != nil || info.IsDir() {
return nil
}
if !isVideo(path) {
return nil
}
found[path] = true
s.indexFile(path, info, "FILE_UPDATED", mediaType, shelfID, emitEvents)
return nil
})
}
func (s *Scanner) watchFilesystem() {
for {
s.runWatchSession()
}
}
func (s *Scanner) runWatchSession() {
watcher, err := fsnotify.NewWatcher()
if err != nil {
log.Printf("fsnotify unavailable: %v", err)
time.Sleep(5 * time.Second)
return
}
defer watcher.Close()
rootType := map[string]string{}
rootShelf := map[string]string{}
watchDirs := 0
addTree := func(dir, kind, shelfID string) {
if dir == "" {
return
}
clean := filepath.Clean(dir)
rootType[clean] = kind
if shelfID != "" {
rootShelf[clean] = shelfID
}
filepath.Walk(dir, func(path string, info os.FileInfo, err error) error {
if err == nil && info.IsDir() {
if err := watcher.Add(path); err != nil {
log.Printf("watch add failed for %s: %v", path, err)
} else {
watchDirs++
}
}
return nil
})
}
s.mu.Lock()
roots := s.cfg.Media.EffectiveRoots()
s.mu.Unlock()
for _, root := range roots {
kind := root.Kind
if kind == "series" {
kind = "episode"
}
addTree(root.Path, kind, root.ShelfID)
}
log.Printf("Filesystem watch active on %d roots (%d directories)", len(roots), watchDirs)
debounce := time.NewTimer(0)
<-debounce.C
pending := make(map[string]struct{})
const settle = 5 * time.Second
for {
select {
case <-s.reloadWatch:
log.Println("Reloading filesystem watchers for new media paths...")
return
case event, ok := <-watcher.Events:
if !ok {
return
}
name := event.Name
if event.Op&fsnotify.Create == fsnotify.Create {
if info, err := os.Stat(name); err == nil && info.IsDir() {
if err := watcher.Add(name); err != nil {
log.Printf("watch add failed for %s: %v", name, err)
}
filepath.Walk(name, func(path string, info os.FileInfo, err error) error {
if err != nil {
return nil
}
if info.IsDir() {
if err := watcher.Add(path); err != nil {
log.Printf("watch add failed for %s: %v", path, err)
}
return nil
}
if isVideo(path) {
pending[path] = struct{}{}
}
return nil
})
}
}
if event.Op&fsnotify.Create == fsnotify.Create ||
event.Op&fsnotify.Write == fsnotify.Write ||
event.Op&fsnotify.Rename == fsnotify.Rename {
if isVideo(name) {
pending[name] = struct{}{}
}
}
if event.Op&fsnotify.Remove == fsnotify.Remove || event.Op&fsnotify.Rename == fsnotify.Rename {
s.handleRemove(name)
}
if !debounce.Stop() {
select {
case <-debounce.C:
default:
}
}
debounce.Reset(settle)
case <-debounce.C:
for path := range pending {
info, err := os.Stat(path)
if err == nil && !info.IsDir() && isVideo(path) {
kind, shelf := s.detectRootMeta(path, rootType, rootShelf)
s.indexFile(path, info, "FILE_ADDED", kind, shelf, true)
}
}
pending = make(map[string]struct{})
case err, ok := <-watcher.Errors:
if !ok {
return
}
log.Printf("watch error: %v", err)
}
}
}
func (s *Scanner) detectRootMeta(path string, rootType, rootShelf map[string]string) (kind, shelfID string) {
clean := filepath.Clean(path)
bestRoot := ""
bestKind := ""
for root, k := range rootType {
if clean == root || strings.HasPrefix(clean, root+string(filepath.Separator)) {
if len(root) >= len(bestRoot) {
bestRoot = root
bestKind = k
}
}
}
if bestRoot != "" {
return bestKind, rootShelf[bestRoot]
}
return detectMediaType(path), ""
}
func (s *Scanner) detectRootType(path string, rootType map[string]string) string {
clean := filepath.Clean(path)
bestRoot := ""
bestKind := ""
for root, kind := range rootType {
if clean == root || strings.HasPrefix(clean, root+string(filepath.Separator)) {
if len(root) >= len(bestRoot) {
bestRoot = root
bestKind = kind
}
}
}
if bestRoot != "" {
return bestKind
}
return detectMediaType(path)
}
func (s *Scanner) indexFile(path string, info os.FileInfo, eventType, preferredType, shelfID string, emitEvents bool) {
mediaType := preferredType
if mediaType == "" {
mediaType = detectMediaType(path)
}
if mediaType == "" {
return
}
if mediaType == "movie" && parseMovieFromPath(path) == nil {
return
}
if mediaType == "episode" && parseSeriesFromPath(path) == nil {
return
}
localFileID := database.FileID(path)
var rev int64
if emitEvents {
rev, _ = s.store.NextRevision()
}
lf := database.LocalFile{
LocalFileID: localFileID,
Path: path,
SizeBytes: info.Size(),
ModifiedAt: info.ModTime(),
MediaType: mediaType,
ReleaseName: filepath.Base(path),
Revision: rev,
}
_ = s.store.UpsertFile(lf)
if !emitEvents {
return
}
info2 := buildMediaInfo(lf)
info2.ShelfID = shelfID
event := LibraryEvent{
Type: eventType,
NodeRevision: rev,
File: info2,
}
s.onEvents([]LibraryEvent{event})
}
// IndexPath forces a library event for a newly imported video file.
func (s *Scanner) IndexPath(path string) {
info, err := os.Stat(path)
if err != nil || info.IsDir() || !isVideo(path) {
return
}
s.mu.Lock()
roots := s.cfg.Media.EffectiveRoots()
s.mu.Unlock()
rootType := map[string]string{}
rootShelf := map[string]string{}
for _, root := range roots {
kind := root.Kind
if kind == "series" {
kind = "episode"
}
clean := filepath.Clean(root.Path)
rootType[clean] = kind
if root.ShelfID != "" {
rootShelf[clean] = root.ShelfID
}
}
kind, shelf := s.detectRootMeta(path, rootType, rootShelf)
s.indexFile(path, info, "FILE_ADDED", kind, shelf, true)
}
func (s *Scanner) handleRemove(path string) {
files, err := s.store.FilesUnderPath(path)
if err != nil || len(files) == 0 {
// Exact file id lookup for leaf files
localFileID := database.FileID(path)
f, getErr := s.store.GetFile(localFileID)
if getErr != nil {
return
}
files = []database.LocalFile{*f}
}
for _, f := range files {
rev, _ := s.store.NextRevision()
event := LibraryEvent{
Type: "FILE_REMOVED",
NodeRevision: rev,
File: fileToInfo(f),
}
_ = s.store.DeleteFile(f.LocalFileID)
s.onEvents([]LibraryEvent{event})
}
}
func (s *Scanner) AllMediaFiles() []MediaFileInfo {
files, _ := s.store.AllFiles()
s.mu.Lock()
roots := s.cfg.Media.EffectiveRoots()
s.mu.Unlock()
result := make([]MediaFileInfo, len(files))
for i, f := range files {
info := buildMediaInfo(f)
info.ShelfID = shelfForPath(f.Path, roots)
result[i] = info
}
return result
}
func shelfForPath(path string, roots []config.MediaRoot) string {
clean := filepath.Clean(path)
best := ""
bestLen := -1
for _, r := range roots {
root := filepath.Clean(r.Path)
if clean == root || strings.HasPrefix(clean, root+string(filepath.Separator)) {
if len(root) > bestLen {
bestLen = len(root)
best = r.ShelfID
}
}
}
return best
}
func isVideo(path string) bool {
ext := strings.ToLower(filepath.Ext(path))
return videoExtensions[ext]
}
func detectMediaType(path string) string {
if parseSeriesFromPath(path) != nil {
return "episode"
}
if parseMovieFromPath(path) != nil {
return "movie"
}
return ""
}
func buildMediaInfo(f database.LocalFile) MediaFileInfo {
info := MediaFileInfo{
LocalFileID: f.LocalFileID,
SizeBytes: f.SizeBytes,
Container: strings.TrimPrefix(filepath.Ext(f.ReleaseName), "."),
ReleaseName: f.ReleaseName,
ModifiedAt: f.ModifiedAt.UTC().Format(time.RFC3339),
MediaType: f.MediaType,
}
if f.MediaType == "movie" {
if p := parseMovieFromPath(f.Path); p != nil {
info.ParsedMovie = p
info.Resolution = strVal(p, "resolution")
info.VideoCodec = strVal(p, "videoCodec")
info.AudioCodec = strVal(p, "audio")
}
} else {
if p := parseSeriesFromPath(f.Path); p != nil {
info.ParsedEpisode = p
info.Resolution = strVal(p, "resolution")
info.VideoCodec = strVal(p, "videoCodec")
info.AudioCodec = strVal(p, "audio")
}
}
return info
}
func fileToInfo(f database.LocalFile) MediaFileInfo {
return buildMediaInfo(f)
}
func strVal(m map[string]interface{}, key string) string {
if v, ok := m[key].(string); ok {
return v
}
return ""
}