mirror of
https://github.com/truewhile/MeBox.git
synced 2026-10-05 13:06:36 +08:00
Import local media before track probing
This commit is contained in:
+114
-16
@@ -74,6 +74,10 @@ type ScannerService struct {
|
|||||||
cloudMediaProbeBackoff map[string]time.Time
|
cloudMediaProbeBackoff map[string]time.Time
|
||||||
cloudMediaProbeWarnMu sync.Mutex
|
cloudMediaProbeWarnMu sync.Mutex
|
||||||
cloudMediaProbeLastWarn time.Time
|
cloudMediaProbeLastWarn time.Time
|
||||||
|
localMediaProbeOnce sync.Once
|
||||||
|
localMediaProbeQueue chan localMediaProbeTask
|
||||||
|
localMediaProbeMu sync.Mutex
|
||||||
|
localMediaProbing map[string]struct{}
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewScannerService is the constructor.
|
// NewScannerService is the constructor.
|
||||||
@@ -96,6 +100,8 @@ func NewScannerService(
|
|||||||
cloudMediaProbeQueue: make(chan cloudMediaProbeTask, 1024),
|
cloudMediaProbeQueue: make(chan cloudMediaProbeTask, 1024),
|
||||||
cloudMediaProbing: make(map[string]struct{}),
|
cloudMediaProbing: make(map[string]struct{}),
|
||||||
cloudMediaProbeBackoff: make(map[string]time.Time),
|
cloudMediaProbeBackoff: make(map[string]time.Time),
|
||||||
|
localMediaProbeQueue: make(chan localMediaProbeTask, 1024),
|
||||||
|
localMediaProbing: make(map[string]struct{}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -106,7 +112,7 @@ func (s *ScannerService) SetStorageConfig(storage *StorageConfigService) {
|
|||||||
s.storage = storage
|
s.storage = storage
|
||||||
if storage != nil && s.probe != nil {
|
if storage != nil && s.probe != nil {
|
||||||
s.cloudMediaProbeOnce.Do(func() {
|
s.cloudMediaProbeOnce.Do(func() {
|
||||||
workers := normalizeFFprobeMaxConcurrent(s.cfg.App.FFprobeMaxConcurrent)
|
workers := s.ffprobeWorkerCount()
|
||||||
for i := 0; i < workers; i++ {
|
for i := 0; i < workers; i++ {
|
||||||
go s.cloudMediaProbeWorker()
|
go s.cloudMediaProbeWorker()
|
||||||
}
|
}
|
||||||
@@ -343,6 +349,10 @@ type cloudMediaProbeTask struct {
|
|||||||
path string
|
path string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type localMediaProbeTask struct {
|
||||||
|
path string
|
||||||
|
}
|
||||||
|
|
||||||
type existingCloudMedia struct {
|
type existingCloudMedia struct {
|
||||||
SizeBytes int64
|
SizeBytes int64
|
||||||
DurationSec int
|
DurationSec int
|
||||||
@@ -448,6 +458,52 @@ func (s *ScannerService) queueCloudMediaProbeWithBudget(typ, ref, path string, b
|
|||||||
return s.queueCloudMediaProbe(typ, ref, path)
|
return s.queueCloudMediaProbe(typ, ref, path)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *ScannerService) localMediaProbeWorker() {
|
||||||
|
for task := range s.localMediaProbeQueue {
|
||||||
|
s.probeLocalMediaAsync(task)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *ScannerService) queueLocalMediaProbe(path string) bool {
|
||||||
|
if s == nil || s.probe == nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
path = strings.TrimSpace(path)
|
||||||
|
if path == "" {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
s.localMediaProbeOnce.Do(func() {
|
||||||
|
workers := s.ffprobeWorkerCount()
|
||||||
|
for i := 0; i < workers; i++ {
|
||||||
|
go s.localMediaProbeWorker()
|
||||||
|
}
|
||||||
|
})
|
||||||
|
s.localMediaProbeMu.Lock()
|
||||||
|
if s.localMediaProbing == nil {
|
||||||
|
s.localMediaProbing = make(map[string]struct{})
|
||||||
|
}
|
||||||
|
if _, ok := s.localMediaProbing[path]; ok {
|
||||||
|
s.localMediaProbeMu.Unlock()
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
s.localMediaProbing[path] = struct{}{}
|
||||||
|
s.localMediaProbeMu.Unlock()
|
||||||
|
|
||||||
|
task := localMediaProbeTask{path: path}
|
||||||
|
select {
|
||||||
|
case s.localMediaProbeQueue <- task:
|
||||||
|
return true
|
||||||
|
default:
|
||||||
|
s.localMediaProbeMu.Lock()
|
||||||
|
delete(s.localMediaProbing, path)
|
||||||
|
s.localMediaProbeMu.Unlock()
|
||||||
|
if s.log != nil {
|
||||||
|
s.log.Debug("local media probe queue full", zap.String("path", path))
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (s *ScannerService) beginCloudScan(ctx context.Context, lib *model.Library, mount CloudMountInfo) (context.Context, func(*ScanResult, error), error) {
|
func (s *ScannerService) beginCloudScan(ctx context.Context, lib *model.Library, mount CloudMountInfo) (context.Context, func(*ScanResult, error), error) {
|
||||||
if s == nil || lib == nil {
|
if s == nil || lib == nil {
|
||||||
return ctx, func(*ScanResult, error) {}, nil
|
return ctx, func(*ScanResult, error) {}, nil
|
||||||
@@ -1448,6 +1504,55 @@ func (s *ScannerService) probeCloudMediaAsync(task cloudMediaProbeTask) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *ScannerService) probeLocalMediaAsync(task localMediaProbeTask) {
|
||||||
|
defer func() {
|
||||||
|
s.localMediaProbeMu.Lock()
|
||||||
|
delete(s.localMediaProbing, task.path)
|
||||||
|
s.localMediaProbeMu.Unlock()
|
||||||
|
}()
|
||||||
|
if s == nil || s.probe == nil || strings.TrimSpace(task.path) == "" {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
probe, err := s.probe.Probe(ctx, task.path)
|
||||||
|
if err != nil {
|
||||||
|
if s.log != nil {
|
||||||
|
s.log.Debug("local media async probe failed", zap.String("path", task.path), zap.Error(err))
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
updates := probeResultUpdates(probe)
|
||||||
|
if len(updates) == 0 {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if err := s.repo.DB.WithContext(ctx).Model(&model.Media{}).Where("path = ?", task.path).Updates(updates).Error; err != nil {
|
||||||
|
if s.log != nil {
|
||||||
|
s.log.Debug("update local media track metadata failed", zap.String("path", task.path), zap.Error(err))
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if s.hub != nil {
|
||||||
|
s.hub.Publish("scan", map[string]any{
|
||||||
|
"path": task.path,
|
||||||
|
"track_probed": true,
|
||||||
|
"duration_sec": probe.DurationSec,
|
||||||
|
"video_codec": probe.VideoCodec,
|
||||||
|
"audio_codec": probe.AudioCodec,
|
||||||
|
"width": probe.Width,
|
||||||
|
"height": probe.Height,
|
||||||
|
"probe_message": "本地媒体轨道元数据已后台补齐",
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *ScannerService) ffprobeWorkerCount() int {
|
||||||
|
if s == nil || s.cfg == nil {
|
||||||
|
return 1
|
||||||
|
}
|
||||||
|
return normalizeFFprobeMaxConcurrent(s.cfg.App.FFprobeMaxConcurrent)
|
||||||
|
}
|
||||||
|
|
||||||
func (s *ScannerService) probeCloudFileMetadata(ctx context.Context, typ, ref string) (*ProbeResult, error) {
|
func (s *ScannerService) probeCloudFileMetadata(ctx context.Context, typ, ref string) (*ProbeResult, error) {
|
||||||
if s == nil || s.probe == nil || s.storage == nil {
|
if s == nil || s.probe == nil || s.storage == nil {
|
||||||
return nil, errors.New("cloud probe unavailable")
|
return nil, errors.New("cloud probe unavailable")
|
||||||
@@ -1648,25 +1753,15 @@ func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, pat
|
|||||||
res.LocalMetadata++
|
res.LocalMetadata++
|
||||||
}
|
}
|
||||||
|
|
||||||
// Best-effort ffprobe; failure does not abort the file.
|
var after func()
|
||||||
if s.probe != nil {
|
if ext != ".strm" && s.probe != nil {
|
||||||
if probe, err := s.probe.Probe(ctx, path); err == nil && probe != nil {
|
after = func() {
|
||||||
m.DurationSec = probe.DurationSec
|
s.queueLocalMediaProbe(path)
|
||||||
m.Width = probe.Width
|
|
||||||
m.Height = probe.Height
|
|
||||||
m.VideoCodec = probe.VideoCodec
|
|
||||||
m.AudioCodec = probe.AudioCodec
|
|
||||||
if probe.Container != "" {
|
|
||||||
m.Container = probe.Container
|
|
||||||
}
|
|
||||||
res.Probed++
|
|
||||||
} else if err != nil {
|
|
||||||
s.log.Debug("ffprobe failed", zap.String("path", path), zap.Error(err))
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if isNewMedia && writeBatch != nil {
|
if isNewMedia && writeBatch != nil {
|
||||||
writeBatch.Add(path, m)
|
writeBatch.AddWithAfter(path, m, after)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if err := s.repo.Media.Upsert(ctx, m); err != nil {
|
if err := s.repo.Media.Upsert(ctx, m); err != nil {
|
||||||
@@ -1674,6 +1769,9 @@ func (s *ScannerService) ingestFile(ctx context.Context, lib *model.Library, pat
|
|||||||
s.log.Warn("upsert media failed", zap.String("path", path), zap.Error(err))
|
s.log.Warn("upsert media failed", zap.String("path", path), zap.Error(err))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
if after != nil {
|
||||||
|
after()
|
||||||
|
}
|
||||||
if isNewMedia {
|
if isNewMedia {
|
||||||
res.Added++
|
res.Added++
|
||||||
} else {
|
} else {
|
||||||
|
|||||||
@@ -155,6 +155,28 @@ func TestScanLibrarySkipsUnchangedExistingLocalMediaWithMissingTrackMetadata(t *
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestScanLibraryImportsNewLocalMediaWithoutSynchronousProbe(t *testing.T) {
|
||||||
|
sc, repos := newScannerTestEnv(t)
|
||||||
|
sc.probe = NewFFprobeService(&config.Config{}, zap.NewNop())
|
||||||
|
root := t.TempDir()
|
||||||
|
lib := model.Library{Name: "Movies", Path: root, Type: "movie", Enabled: true}
|
||||||
|
if err := repos.Library.Create(t.Context(), &lib); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
file := filepath.Join(root, "New Movie (2026).mkv")
|
||||||
|
if err := os.WriteFile(file, []byte("new-file"), 0o644); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
res, err := sc.ScanLibrary(t.Context(), lib.ID)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("scan: %v", err)
|
||||||
|
}
|
||||||
|
if res.Added != 1 || res.Probed != 0 {
|
||||||
|
t.Fatalf("scan result = %#v, want fast import without synchronous ffprobe", res)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestScanLibraryReportsPerFileUpsertErrors(t *testing.T) {
|
func TestScanLibraryReportsPerFileUpsertErrors(t *testing.T) {
|
||||||
sc, repos := newScannerTestEnv(t)
|
sc, repos := newScannerTestEnv(t)
|
||||||
root := t.TempDir()
|
root := t.TempDir()
|
||||||
|
|||||||
Reference in New Issue
Block a user