mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-03 12:26:36 +08:00
Preserve media metadata across STRM path renames
This commit is contained in:
@@ -7,6 +7,7 @@ import (
|
||||
"go.uber.org/zap"
|
||||
|
||||
"github.com/truewhile/MeBox/internal/model"
|
||||
"github.com/truewhile/MeBox/internal/repository"
|
||||
)
|
||||
|
||||
type localMediaWriteBatch struct {
|
||||
@@ -18,9 +19,10 @@ type localMediaWriteBatch struct {
|
||||
}
|
||||
|
||||
type localMediaWriteItem struct {
|
||||
path string
|
||||
media *model.Media
|
||||
after func()
|
||||
path string
|
||||
media *model.Media
|
||||
aliases []string
|
||||
after func()
|
||||
}
|
||||
|
||||
func newLocalMediaWriteBatch(scanner *ScannerService, ctx context.Context, res *ScanResult, limit int) *localMediaWriteBatch {
|
||||
@@ -41,7 +43,12 @@ func (b *localMediaWriteBatch) AddWithAfter(path string, media *model.Media, aft
|
||||
if media.ScrapeStatus == "" {
|
||||
media.ScrapeStatus = "pending"
|
||||
}
|
||||
b.items = append(b.items, localMediaWriteItem{path: path, media: media, after: after})
|
||||
b.items = append(b.items, localMediaWriteItem{
|
||||
path: path,
|
||||
media: media,
|
||||
aliases: missingSTRMPathAliases(path),
|
||||
after: after,
|
||||
})
|
||||
if len(b.items) >= b.limit {
|
||||
b.Flush()
|
||||
}
|
||||
@@ -53,39 +60,35 @@ func (b *localMediaWriteBatch) Flush() {
|
||||
}
|
||||
items := b.items
|
||||
b.items = nil
|
||||
media := make([]model.Media, 0, len(items))
|
||||
for _, item := range items {
|
||||
if item.media != nil {
|
||||
media = append(media, *item.media)
|
||||
}
|
||||
}
|
||||
if len(media) == 0 {
|
||||
return
|
||||
}
|
||||
existingPaths := b.existingPaths(items)
|
||||
upsertItems := make([]*model.Media, 0, len(items))
|
||||
upsertAfter := make([]func(), 0, len(items))
|
||||
|
||||
existingPaths, lookupOK := b.existingPathOrAliasSet(items)
|
||||
upsertItems := make([]localMediaWriteItem, 0, len(items))
|
||||
createItems := make([]localMediaWriteItem, 0, len(items))
|
||||
createMedia := make([]model.Media, 0, len(items))
|
||||
for _, item := range items {
|
||||
if item.media == nil {
|
||||
continue
|
||||
}
|
||||
if existingPaths[filepath.Clean(item.media.Path)] {
|
||||
// 已存在行:攒起来在一个事务里逐条 upsert(一批一次提交)。
|
||||
after := item.after
|
||||
upsertItems = append(upsertItems, item.media)
|
||||
upsertAfter = append(upsertAfter, after)
|
||||
// If the alias lookup failed, route through UpsertWithAliases instead of
|
||||
// direct-create. That is slower but cannot create a duplicate row.
|
||||
if !lookupOK || mediaPathOrAliasExists(item.media.Path, item.aliases, existingPaths) {
|
||||
upsertItems = append(upsertItems, item)
|
||||
continue
|
||||
}
|
||||
createItems = append(createItems, item)
|
||||
createMedia = append(createMedia, *item.media)
|
||||
}
|
||||
b.flushUpserts(items, upsertItems, upsertAfter)
|
||||
if len(createMedia) == 0 {
|
||||
|
||||
b.flushUpserts(upsertItems)
|
||||
if len(createItems) == 0 {
|
||||
b.publish()
|
||||
return
|
||||
}
|
||||
|
||||
createMedia := make([]model.Media, 0, len(createItems))
|
||||
for _, item := range createItems {
|
||||
if item.media != nil {
|
||||
createMedia = append(createMedia, *item.media)
|
||||
}
|
||||
}
|
||||
if err := b.scanner.repo.DB.WithContext(b.ctx).CreateInBatches(&createMedia, b.limit).Error; err == nil {
|
||||
b.res.Added += len(createMedia)
|
||||
for _, item := range createItems {
|
||||
@@ -96,12 +99,13 @@ func (b *localMediaWriteBatch) Flush() {
|
||||
b.publish()
|
||||
return
|
||||
}
|
||||
|
||||
for _, item := range createItems {
|
||||
if item.media == nil {
|
||||
continue
|
||||
}
|
||||
wasExisting := b.mediaPathExists(item.media.Path)
|
||||
if err := b.scanner.repo.Media.Upsert(b.ctx, item.media); err != nil {
|
||||
if err := b.scanner.repo.Media.UpsertWithAliases(b.ctx, item.media, item.aliases); err != nil {
|
||||
addScanError(b.res, item.path, err)
|
||||
b.scanner.log.Warn("upsert media failed", zap.String("path", item.path), zap.Error(err))
|
||||
continue
|
||||
@@ -118,47 +122,90 @@ func (b *localMediaWriteBatch) Flush() {
|
||||
b.publish()
|
||||
}
|
||||
|
||||
func (b *localMediaWriteBatch) existingPaths(items []localMediaWriteItem) map[string]bool {
|
||||
// existingPathOrAliasSet loads exact and alias paths in bounded query chunks. Aliases
|
||||
// have already been filtered against the filesystem by AddWithAfter, so an
|
||||
// existing keep_ext sibling is never treated as a rename.
|
||||
func (b *localMediaWriteBatch) existingPathOrAliasSet(items []localMediaWriteItem) (map[string]bool, bool) {
|
||||
out := map[string]bool{}
|
||||
if b == nil || b.scanner == nil || b.scanner.repo == nil || b.scanner.repo.DB == nil || len(items) == 0 {
|
||||
return out
|
||||
return out, false
|
||||
}
|
||||
seen := make(map[string]struct{}, len(items)*2)
|
||||
paths := make([]string, 0, len(items)*2)
|
||||
add := func(path string) {
|
||||
path = filepath.Clean(path)
|
||||
if path == "" || path == "." {
|
||||
return
|
||||
}
|
||||
if _, ok := seen[path]; ok {
|
||||
return
|
||||
}
|
||||
seen[path] = struct{}{}
|
||||
paths = append(paths, path)
|
||||
}
|
||||
paths := make([]string, 0, len(items))
|
||||
for _, item := range items {
|
||||
if item.media == nil || item.media.Path == "" {
|
||||
continue
|
||||
}
|
||||
paths = append(paths, item.media.Path)
|
||||
add(item.media.Path)
|
||||
for _, alias := range item.aliases {
|
||||
add(alias)
|
||||
}
|
||||
}
|
||||
if len(paths) == 0 {
|
||||
return out
|
||||
return out, true
|
||||
}
|
||||
var rows []string
|
||||
if err := b.scanner.repo.DB.WithContext(b.ctx).
|
||||
Unscoped().
|
||||
Model(&model.Media{}).
|
||||
Where("path IN ?", paths).
|
||||
Pluck("path", &rows).Error; err != nil {
|
||||
b.scanner.log.Debug("load existing media paths for scan batch failed", zap.Error(err))
|
||||
return out
|
||||
// Keep SQL variables below the legacy SQLite limit (999). A batch can
|
||||
// contain 100 STRM rows and every row has many sibling aliases.
|
||||
const pathLookupChunk = 400
|
||||
for start := 0; start < len(paths); start += pathLookupChunk {
|
||||
end := start + pathLookupChunk
|
||||
if end > len(paths) {
|
||||
end = len(paths)
|
||||
}
|
||||
var rows []string
|
||||
if err := b.scanner.repo.DB.WithContext(b.ctx).
|
||||
Unscoped().
|
||||
Model(&model.Media{}).
|
||||
Where("path IN ?", paths[start:end]).
|
||||
Pluck("path", &rows).Error; err != nil {
|
||||
b.scanner.log.Debug("load existing media paths for scan batch failed", zap.Error(err))
|
||||
return nil, false
|
||||
}
|
||||
for _, path := range rows {
|
||||
out[filepath.Clean(path)] = true
|
||||
}
|
||||
}
|
||||
for _, path := range rows {
|
||||
out[filepath.Clean(path)] = true
|
||||
}
|
||||
return out
|
||||
return out, true
|
||||
}
|
||||
|
||||
// flushUpserts 把已存在行的 upsert 攒成一个事务(一次提交/一组 fsync)。
|
||||
// 整批失败(如单条数据触发约束)时退回逐条 Upsert,只丢真正坏的那几条。
|
||||
func (b *localMediaWriteBatch) flushUpserts(allItems []localMediaWriteItem, upsertItems []*model.Media, upsertAfter []func()) {
|
||||
func mediaPathOrAliasExists(path string, aliases []string, existing map[string]bool) bool {
|
||||
if existing[filepath.Clean(path)] {
|
||||
return true
|
||||
}
|
||||
for _, alias := range aliases {
|
||||
if existing[filepath.Clean(alias)] {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
// flushUpserts 把已存在或可迁移路径的行攒成一个事务(一次提交/一组 fsync)。
|
||||
// 整批失败(如单条数据触发约束)时退回逐条 upsert,只丢真正坏的那几条。
|
||||
func (b *localMediaWriteBatch) flushUpserts(upsertItems []localMediaWriteItem) {
|
||||
if len(upsertItems) == 0 {
|
||||
return
|
||||
}
|
||||
if err := b.scanner.repo.Media.UpsertBatch(b.ctx, upsertItems); err == nil {
|
||||
batchItems := make([]repository.MediaUpsertItem, 0, len(upsertItems))
|
||||
for _, item := range upsertItems {
|
||||
batchItems = append(batchItems, repository.MediaUpsertItem{Media: item.media, AliasPaths: item.aliases})
|
||||
}
|
||||
if err := b.scanner.repo.Media.UpsertBatchWithAliases(b.ctx, batchItems); err == nil {
|
||||
b.res.Updated += len(upsertItems)
|
||||
for _, after := range upsertAfter {
|
||||
if after != nil {
|
||||
after()
|
||||
for _, item := range upsertItems {
|
||||
if item.after != nil {
|
||||
item.after()
|
||||
}
|
||||
}
|
||||
return
|
||||
@@ -166,13 +213,7 @@ func (b *localMediaWriteBatch) flushUpserts(allItems []localMediaWriteItem, upse
|
||||
b.scanner.log.Warn("batch upsert failed; falling back to per-item upsert",
|
||||
zap.Int("items", len(upsertItems)))
|
||||
}
|
||||
// 兜底:按原始顺序找回每个条目的 path/after(两个切片同序但可能含 nil)。
|
||||
idx := 0
|
||||
for _, item := range allItems {
|
||||
if item.media == nil || idx >= len(upsertItems) || upsertItems[idx] != item.media {
|
||||
continue
|
||||
}
|
||||
idx++
|
||||
for _, item := range upsertItems {
|
||||
b.upsertExistingItem(item)
|
||||
}
|
||||
}
|
||||
@@ -181,7 +222,7 @@ func (b *localMediaWriteBatch) upsertExistingItem(item localMediaWriteItem) {
|
||||
if item.media == nil {
|
||||
return
|
||||
}
|
||||
if err := b.scanner.repo.Media.Upsert(b.ctx, item.media); err != nil {
|
||||
if err := b.scanner.repo.Media.UpsertWithAliases(b.ctx, item.media, item.aliases); err != nil {
|
||||
addScanError(b.res, item.path, err)
|
||||
b.scanner.log.Warn("upsert media failed", zap.String("path", item.path), zap.Error(err))
|
||||
return
|
||||
|
||||
Reference in New Issue
Block a user