Files
MusicServer/internal/backend/sync.go
T

565 lines
19 KiB
Go

package backend
import (
"crypto/md5"
"encoding/hex"
"errors"
"fmt"
"io"
"io/fs"
"music-server/internal/db/repository"
"music-server/internal/logging"
"os"
"sort"
"strconv"
"strings"
"sync"
"time"
"github.com/panjf2000/ants/v2"
"github.com/MShekow/directory-checksum/directory_checksum"
"github.com/spf13/afero"
"go.uber.org/zap"
)
var Syncing = false
var foldersSynced float32
var numberOfFoldersToSync float32
var start time.Time
var totalTime time.Duration
var timeSpent time.Duration
var allSoundtracks []repository.Soundtrack
var soundtracksBeforeSync []repository.Soundtrack
var soundtracksAfterSync []repository.Soundtrack
var soundtracksAdded []string
var soundtracksReAdded []string
var soundtracksChangedTitle map[string]string
var soundtracksChangedContent []string
var soundtracksRemoved []string
var catchedErrors []string
type brokenSong struct {
SoundtrackID int32
Path string
}
var brokenSongs []brokenSong
var pool *ants.Pool
var poolSong *ants.Pool
type SyncResponse struct {
SoundtracksAdded []string `json:"soundtracks_added"`
SoundtracksReAdded []string `json:"soundtracks_re_added"`
SoundtracksChangedTitle map[string]string `json:"soundtracks_changed_title"`
SoundtracksChangedContent []string `json:"soundtracks_changed_content"`
SoundtracksRemoved []string `json:"soundtracks_removed"`
CatchedErrors []string `json:"catched_errors"`
TotalTime string `json:"total_time"`
}
type ProgressResponse struct {
Progress string `json:"progress"`
TimeSpent string `json:"time_spent"`
}
type SoundtrackStatus int
const (
NotChanged SoundtrackStatus = iota
TitleChanged
SoundtrackChanged
NewSoundtrack
)
var statusName = map[SoundtrackStatus]string{
NotChanged: "Not changed",
TitleChanged: "Title changed",
SoundtrackChanged: "Soundtrack changed",
NewSoundtrack: "New soundtrack",
}
func (ss SoundtrackStatus) String() string {
return statusName[ss]
}
func ResetDB() {
repo.ClearSongs(BackendCtx())
repo.ClearSoundtracks(BackendCtx())
}
func SyncProgress() ProgressResponse {
progress := int((foldersSynced / numberOfFoldersToSync) * 100)
currentTime := time.Now()
timeSpent = currentTime.Sub(start)
out := time.Time{}.Add(timeSpent)
logging.GetLogger().Debug("Sync progress",
zap.Int("progress_percent", progress),
zap.Int("folders_synced", int(foldersSynced)),
zap.Int("total_folders", int(numberOfFoldersToSync)),
zap.String("time_spent", out.Format("15:04:05.00000")))
return ProgressResponse{
Progress: fmt.Sprintf("%v", progress),
TimeSpent: out.Format("15:04:05"),
}
}
func SyncResult() SyncResponse {
logging.GetLogger().Info("Sync completed",
zap.Int("soundtracks_before", len(soundtracksBeforeSync)),
zap.Int("soundtracks_after", len(soundtracksAfterSync)))
if len(soundtracksAdded) > 0 {
logging.GetLogger().Debug("Soundtracks added", zap.Strings("soundtracks", soundtracksAdded))
}
if len(soundtracksReAdded) > 0 {
logging.GetLogger().Debug("Soundtracks readded", zap.Strings("soundtracks", soundtracksReAdded))
}
if len(soundtracksChangedTitle) > 0 {
logging.GetLogger().Debug("Soundtracks with changed title", zap.Any("changes", soundtracksChangedTitle))
}
if len(soundtracksChangedContent) > 0 {
logging.GetLogger().Debug("Soundtracks with changed content", zap.Strings("soundtracks", soundtracksChangedContent))
}
var soundtracksRemovedTemp []string
for _, beforeSoundtrack := range soundtracksBeforeSync {
var found = false
for _, afterSoundtrack := range soundtracksAfterSync {
if beforeSoundtrack.SoundtrackName == afterSoundtrack.SoundtrackName {
found = true
break
}
}
if !found {
soundtracksRemovedTemp = append(soundtracksRemovedTemp, beforeSoundtrack.SoundtrackName)
}
}
for _, soundtrack := range soundtracksRemovedTemp {
var found bool = false
for key := range soundtracksChangedTitle {
if soundtrack == key {
found = true
break
}
}
if !found {
soundtracksRemoved = append(soundtracksRemoved, soundtrack)
}
}
if len(soundtracksRemoved) > 0 {
logging.GetLogger().Debug("Soundtracks removed", zap.Strings("soundtracks", soundtracksRemoved))
}
if len(catchedErrors) > 0 {
logging.GetLogger().Error("Errors caught during sync", zap.Strings("errors", catchedErrors))
}
out := time.Time{}.Add(totalTime)
logging.GetLogger().Info("Sync completed", zap.String("total_time", out.Format("15:04:05.00000")))
return SyncResponse{
SoundtracksAdded: soundtracksAdded,
SoundtracksReAdded: soundtracksReAdded,
SoundtracksChangedTitle: soundtracksChangedTitle,
SoundtracksChangedContent: soundtracksChangedContent,
SoundtracksRemoved: soundtracksRemoved,
CatchedErrors: catchedErrors,
TotalTime: out.Format("15:04:05"),
}
}
func SyncSoundtracksFull() {
syncSoundtracks(true)
Reset()
}
func SyncSoundtracksOnlyChanges() {
syncSoundtracks(false)
Reset()
}
func syncSoundtracks(full bool) {
musicPath := os.Getenv("MUSIC_PATH")
fmt.Printf("dir: %s\n", musicPath)
logging.GetLogger().Debug("Folder to sync", zap.String("MUSIC_PATH", musicPath))
if !strings.HasSuffix(musicPath, "/") {
musicPath += "/"
}
var syncWg sync.WaitGroup
initRepo()
start = time.Now()
foldersToSkip := []string{".sync", "characters", "dist", "old"}
logging.GetLogger().Debug("Folders to skip during sync", zap.Strings("folders", foldersToSkip))
var err error
soundtracksAdded = nil
soundtracksReAdded = nil
soundtracksChangedTitle = nil
soundtracksChangedContent = nil
soundtracksRemoved = nil
catchedErrors = nil
brokenSongs = nil
soundtracksBeforeSync, err = repo.FindAllSoundtracks(BackendCtx())
handleError("FindAllSoundtracks Before", err, "")
logging.GetLogger().Info("Starting sync", zap.Int("soundtracks_before", len(soundtracksBeforeSync)))
allSoundtracks, err = repo.GetAllSoundtracksIncludingDeleted(BackendCtx())
handleError("GetAllSoundtracksIncludingDeleted", err, "")
err = repo.SetSoundtrackDeletionDate(BackendCtx())
handleError("SetSoundtrackDeletionDate", err, "")
directories, err := os.ReadDir(musicPath)
if err != nil {
logging.GetLogger().Fatal("Failed to read music directory", zap.String("path", musicPath), zap.String("error", err.Error()))
}
pool, _ = ants.NewPool(10, ants.WithPreAlloc(true))
poolSong, _ = ants.NewPool(10, ants.WithPreAlloc(true))
defer pool.Release()
defer poolSong.Release()
foldersSynced = 0
numberOfFoldersToSync = float32(len(directories))
syncWg.Add(int(numberOfFoldersToSync))
for _, dir := range directories {
pool.Submit(func() {
defer syncWg.Done()
syncSoundtrack(dir, foldersToSkip, musicPath, full)
})
}
syncWg.Wait()
checkBrokenSongs()
soundtracksAfterSync, err = repo.FindAllSoundtracks(BackendCtx())
handleError("FindAllSoundtracks After", err, "")
finished := time.Now()
totalTime = finished.Sub(start)
out := time.Time{}.Add(totalTime)
logging.GetLogger().Info("Sync completed", zap.Duration("total_time", totalTime), zap.String("formatted_time", out.Format("15:04:05.00000")))
Syncing = false
}
func checkBrokenSongs() {
allSongs, err := repo.FetchAllSongs(BackendCtx())
handleError("FetchAllSongs", err, "")
var brokenWg sync.WaitGroup
poolBroken, _ := ants.NewPool(200, ants.WithPreAlloc(true))
defer poolBroken.Release()
brokenWg.Add(len(allSongs))
for _, song := range allSongs {
poolBroken.Submit(func() {
defer brokenWg.Done()
checkBrokenSong(song)
})
}
brokenWg.Wait()
for _, bs := range brokenSongs {
err = repo.RemoveBrokenSong(BackendCtx(), repository.RemoveBrokenSongParams{SoundtrackID: bs.SoundtrackID, Path: bs.Path})
handleError("RemoveBrokenSong", err, "")
}
}
func checkBrokenSong(song repository.Song) {
//Check if file exists and open
openFile, err := os.Open(song.Path)
if err != nil {
//File not found
brokenSongs = append(brokenSongs, brokenSong{SoundtrackID: song.SoundtrackID, Path: song.Path})
logging.GetLogger().Warn("Broken song found", zap.String("path", song.Path))
} else {
err = openFile.Close()
if err != nil {
logging.GetLogger().Error("Failed to close file", zap.String("path", song.Path), zap.String("error", err.Error()))
}
}
}
func syncSoundtrack(file os.DirEntry, foldersToSkip []string, baseDir string, full bool) {
if file.IsDir() && !contains(foldersToSkip, file.Name()) {
logging.GetLogger().Debug("Syncing soundtrack", zap.String("soundtrack", file.Name()))
soundtrackDir := baseDir + file.Name() + "/"
dirHash := getHashForDir(soundtrackDir)
var status SoundtrackStatus = NewSoundtrack
var oldSoundtrack repository.Soundtrack
var id int32 = -1
//fmt.Printf("Soundtracks before: %d\n", len(soundtracksBeforeSync))
for _, currentSoundtrack := range allSoundtracks {
oldSoundtrack = currentSoundtrack
//fmt.Printf("%s | %s\n", oldSoundtrack.SoundtrackName, oldSoundtrack.Hash)
if oldSoundtrack.SoundtrackName == file.Name() && oldSoundtrack.Hash == dirHash {
status = NotChanged
id = oldSoundtrack.ID
//fmt.Printf("Soundtrack not changed\n")
break
} else if oldSoundtrack.SoundtrackName == file.Name() && oldSoundtrack.Hash != dirHash {
status = SoundtrackChanged
id = oldSoundtrack.ID
//fmt.Printf("Soundtrack changed\n")
break
} else if oldSoundtrack.SoundtrackName != file.Name() && oldSoundtrack.Hash == dirHash {
status = TitleChanged
id = oldSoundtrack.ID
//fmt.Printf("SoundtrackName changed\n")
break
}
}
if full && status != NewSoundtrack {
status = TitleChanged
}
entries, err := os.ReadDir(soundtrackDir)
if err != nil {
logging.GetLogger().Error("Failed to read soundtrack directory", zap.String("path", soundtrackDir), zap.String("error", err.Error()))
}
switch status {
case NewSoundtrack:
id = insertSoundtrack(file.Name(), soundtrackDir, dirHash)
logging.GetLogger().Debug("New soundtrack detected",
zap.Int32("id", id),
zap.String("soundtrack", file.Name()),
zap.String("hash", dirHash),
zap.String("status", status.String()))
soundtracksAdded = append(soundtracksAdded, file.Name())
checkSongs(entries, soundtrackDir, id)
case SoundtrackChanged:
logging.GetLogger().Debug("Soundtrack changed",
zap.Int32("id", id),
zap.String("soundtrack", file.Name()),
zap.String("hash", dirHash),
zap.String("status", status.String()))
err = repo.UpdateSoundtrackHash(BackendCtx(), repository.UpdateSoundtrackHashParams{Hash: dirHash, ID: id})
handleError("UpdateSoundtrackHash", err, "")
soundtracksChangedContent = append(soundtracksChangedContent, file.Name())
checkSongs(entries, soundtrackDir, id)
case TitleChanged:
logging.GetLogger().Debug("Soundtrack title changed",
zap.Int32("id", id),
zap.String("oldName", oldSoundtrack.SoundtrackName),
zap.String("newName", file.Name()),
zap.String("hash", dirHash),
zap.String("status", status.String()))
err = repo.UpdateSoundtrackName(BackendCtx(), repository.UpdateSoundtrackNameParams{Name: file.Name(), Path: soundtrackDir, ID: id})
handleError("UpdateSoundtrackName", err, "")
checkSongs(entries, soundtrackDir, id)
if soundtracksChangedTitle == nil {
soundtracksChangedTitle = make(map[string]string)
}
soundtracksChangedTitle[oldSoundtrack.SoundtrackName] = file.Name()
case NotChanged:
var found bool = false
for _, beforeSoundtrack := range soundtracksBeforeSync {
if dirHash == beforeSoundtrack.Hash {
found = true
logging.GetLogger().Debug("Soundtrack not changed",
zap.Int32("id", id),
zap.String("newName", file.Name()),
zap.String("hash", dirHash),
zap.String("status", status.String()))
}
}
if !found {
checkSongs(entries, soundtrackDir, id)
soundtracksReAdded = append(soundtracksReAdded, file.Name())
logging.GetLogger().Debug("Soundtrack added again",
zap.Int32("id", id),
zap.String("newName", file.Name()),
zap.String("hash", dirHash),
zap.String("status", status.String()))
}
}
logging.GetLogger().Debug("Soundtrack sync status",
zap.Int32("id", id),
zap.String("soundtrack", file.Name()),
zap.String("hash", dirHash),
zap.String("status", status.String()))
err = repo.RemoveSoundtrackDeletionDate(BackendCtx(), id)
handleError("RemoveSoundtrackDeletionDate", err, "")
}
foldersSynced++
logging.GetLogger().Debug("Sync progress",
zap.Int("folders_synced", int(foldersSynced)),
zap.Int("total_folders", int(numberOfFoldersToSync)),
zap.Int("percent", int((foldersSynced/numberOfFoldersToSync)*100)))
}
func insertSoundtrack(name string, path string, hash string) int32 {
var duplicateError = errors.New("ERROR: duplicate key value violates unique")
id, err := repo.InsertSoundtrack(BackendCtx(), repository.InsertSoundtrackParams{SoundtrackName: name, Path: path, Hash: hash})
handleError("InsertSoundtrack", err, "")
if err != nil {
logging.GetLogger().Warn("ID collision detected, resetting sequence")
if strings.HasPrefix(err.Error(), duplicateError.Error()) {
logging.GetLogger().Debug("Resetting soundtrack ID sequence")
_, err = repo.ResetSoundtrackIdSeq(BackendCtx())
handleError("ResetSoundtrackIdSeq", err, "")
id = insertSoundtrack(name, path, hash)
}
}
return id
}
func checkSongs(entries []os.DirEntry, soundtrackDir string, id int32) int32 {
//hasher := md5.New()
var numberOfSongs int32
numberOfFiles := len(entries)
var songWg sync.WaitGroup
songWg.Add(numberOfFiles)
for _, entry := range entries {
poolSong.Submit(func() {
defer songWg.Done()
if checkSong(entry, soundtrackDir, id) {
numberOfSongs++
}
})
}
songWg.Wait()
return numberOfSongs
}
func checkSong(entry os.DirEntry, soundtrackDir string, id int32) bool {
fileInfo, err := entry.Info()
if err != nil {
logging.GetLogger().Error("Failed to get file info", zap.String("filename", entry.Name()), zap.String("error", err.Error()))
return false
}
if isSong(fileInfo) {
path := soundtrackDir + entry.Name()
songHash := getHashForFile(path)
//numberOfSongs++
fileName := entry.Name()
songName, _ := strings.CutSuffix(fileName, ".mp3")
song, err := repo.GetSongWithHash(BackendCtx(), songHash)
handleError("GetSongWithHash", err, fmt.Sprintf("SoundtrackID: %d | Path: %s | SongName: %s | SongHash: %s", id, path, entry.Name(), songHash))
if err == nil {
if song.SongName == songName && song.Path == path {
return false
}
}
logging.GetLogger().Debug("Song changed",
zap.Int32("soundtrack_id", id),
zap.String("path", path),
zap.String("song_name", songName),
zap.String("song_hash", songHash))
count, err := repo.CheckSongWithHash(BackendCtx(), songHash)
handleError("CheckSongWithHash", err, fmt.Sprintf("SoundtrackID: %d | Path: %s | SongName: %s | SongHash: %s\n", id, path, entry.Name(), songHash))
if err != nil {
count2, err := repo.CheckSong(BackendCtx(), repository.CheckSongParams{SoundtrackID: id, Path: path})
handleError("CheckSong", err, fmt.Sprintf("SoundtrackID: %d | Path: %s | SongName: %s | SongHash: %s\n", id, path, entry.Name(), songHash))
if count2 > 0 {
err = repo.AddHashToSong(BackendCtx(), repository.AddHashToSongParams{Hash: songHash, SoundtrackID: id, Path: path})
handleError("AddHashToSong", err, fmt.Sprintf("SoundtrackID: %d | Path: %s | SongName: %s | SongHash: %s", id, path, entry.Name(), songHash))
count, err = repo.CheckSongWithHash(BackendCtx(), songHash)
handleError("CheckSongWithHash 2", err, fmt.Sprintf("SoundtrackID: %d | Path: %s | SongName: %s | SongHash: %s", id, path, entry.Name(), songHash))
}
}
//count, _ := repo.CheckSong(ctx, path)
if count > 0 {
err = repo.UpdateSong(BackendCtx(), repository.UpdateSongParams{SongName: songName, FileName: &fileName, Path: path, Hash: songHash})
handleError("UpdateSong", err, fmt.Sprintf("SoundtrackID: %d | Path: %s | SongName: %s | SongHash: %s", id, path, entry.Name(), songHash))
} else {
count2, err := repo.CheckSong(BackendCtx(), repository.CheckSongParams{SoundtrackID: id, Path: path})
handleError("CheckSong", err, fmt.Sprintf("SoundtrackID: %d | Path: %s | SongName: %s | SongHash: %s", id, path, entry.Name(), songHash))
if count2 > 0 {
err = repo.AddHashToSong(BackendCtx(), repository.AddHashToSongParams{Hash: songHash, SoundtrackID: id, Path: path})
handleError("AddHashToSong", err, fmt.Sprintf("SoundtrackID: %d | Path: %s | SongName: %s | SongHash: %s", id, path, entry.Name(), songHash))
} else {
err = repo.AddSong(BackendCtx(), repository.AddSongParams{SoundtrackID: id, SongName: songName, Path: path, FileName: &fileName, Hash: songHash})
handleError("AddSong", err, fmt.Sprintf("SoundtrackID: %d | Path: %s | SongName: %s | SongHash: %s", id, path, entry.Name(), songHash))
}
}
return true
} else if isCoverImage(fileInfo) {
//TODO: Later add cover art image here in db
}
return false
}
func handleError(funcName string, err error, msg string) {
var compareError = errors.New("no rows in result set")
if err != nil {
if compareError.Error() != err.Error() {
logging.GetLogger().Error("Database error",
zap.String("function", funcName),
zap.String("error", err.Error()))
if msg != "" {
logging.GetLogger().Debug("Error context", zap.String("message", msg))
catchedErrors = append(catchedErrors, fmt.Sprintf("Func: %s\nError message: %s\nDebug message: %s", funcName, err, msg))
} else {
catchedErrors = append(catchedErrors, fmt.Sprintf("Func: %s\nError message: %s", funcName, err))
}
}
}
}
func getHashForDir(soundtrackDir string) string {
directory, _ := directory_checksum.ScanDirectory(soundtrackDir, afero.NewOsFs())
hash, _ := directory.ComputeDirectoryChecksums()
return hash
}
func getHashForFile(path string) string {
hasher := md5.New()
readFile, err := os.Open(path)
if err != nil {
logging.GetLogger().Fatal("Failed to open file for hashing", zap.String("path", path), zap.String("error", err.Error()))
}
defer readFile.Close()
hasher.Reset()
_, err = io.Copy(hasher, readFile)
if err != nil {
logging.GetLogger().Fatal("Failed to hash file", zap.String("path", path), zap.String("error", err.Error()))
}
return hex.EncodeToString(hasher.Sum(nil))
}
func getIdFromFile(file os.FileInfo) int32 {
name := file.Name()
if !file.IsDir() && strings.HasSuffix(name, ".id") {
name = strings.Replace(name, ".id", "", 1)
name = strings.Replace(name, ".", "", 1)
i, _ := strconv.Atoi(name)
return int32(i)
}
return -1
}
func isSong(entry fs.FileInfo) bool {
return !entry.IsDir() && strings.HasSuffix(entry.Name(), ".mp3")
}
func isCoverImage(entry fs.FileInfo) bool {
return !entry.IsDir() && strings.Contains(entry.Name(), "cover") &&
(strings.HasSuffix(entry.Name(), ".jpg") || strings.HasSuffix(entry.Name(), ".png"))
}
func contains(s []string, searchTerm string) bool {
i := sort.SearchStrings(s, searchTerm)
return i < len(s) && s[i] == searchTerm
}