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 "" }