feat(09-01): add scan control fields and per-scan cancellable context
- Add scanActive, scanCancel, scanPaused, scanPauseCh fields to Library struct - Create scan_control.go with CancelScan, PauseScan, ResumeScan, IsScanActive, IsScanPaused - Thread per-scan scanCtx through walk and worker pipeline - Add waitIfPaused checkpoint before each worker extraction - Skip orphan cleanup and variant generation on cancelled scan - Emit LibraryScanCancelled instead of LibraryScanComplete when cancelled
This commit is contained in:
+111
-59
@@ -84,6 +84,12 @@ type Library struct {
|
|||||||
conf *Config
|
conf *Config
|
||||||
db *database.DB
|
db *database.DB
|
||||||
rescanHooks RescanHooks
|
rescanHooks RescanHooks
|
||||||
|
|
||||||
|
// Scan control fields — protected by mu.
|
||||||
|
scanActive bool
|
||||||
|
scanCancel context.CancelFunc
|
||||||
|
scanPaused bool
|
||||||
|
scanPauseCh chan struct{}
|
||||||
}
|
}
|
||||||
|
|
||||||
// SetRescanHooks provides optional hooks for cross-cutting
|
// SetRescanHooks provides optional hooks for cross-cutting
|
||||||
@@ -176,6 +182,32 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
metrics := newScanMetrics()
|
metrics := newScanMetrics()
|
||||||
scanStart := time.Now()
|
scanStart := time.Now()
|
||||||
|
|
||||||
|
scanCtx, scanCancel := context.WithCancel(l.ctx)
|
||||||
|
defer scanCancel()
|
||||||
|
|
||||||
|
l.mu.Lock()
|
||||||
|
l.scanCancel = scanCancel
|
||||||
|
l.scanActive = true
|
||||||
|
l.scanPaused = false
|
||||||
|
l.scanPauseCh = nil
|
||||||
|
l.mu.Unlock()
|
||||||
|
|
||||||
|
defer func() {
|
||||||
|
l.mu.Lock()
|
||||||
|
l.scanCancel = nil
|
||||||
|
l.scanActive = false
|
||||||
|
// If still paused, unpause so no dangling channel.
|
||||||
|
if l.scanPaused {
|
||||||
|
l.scanPaused = false
|
||||||
|
if l.scanPauseCh != nil {
|
||||||
|
close(l.scanPauseCh)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
l.scanPauseCh = nil
|
||||||
|
l.mu.Unlock()
|
||||||
|
}()
|
||||||
|
|
||||||
if len(l.conf.DirectoryPath) == 0 {
|
if len(l.conf.DirectoryPath) == 0 {
|
||||||
return metrics, errLibraryDirNotConfigured
|
return metrics, errLibraryDirNotConfigured
|
||||||
}
|
}
|
||||||
@@ -294,8 +326,8 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
needsUpdate: true,
|
needsUpdate: true,
|
||||||
existingLength: audioFile.LengthMilliseconds,
|
existingLength: audioFile.LengthMilliseconds,
|
||||||
}:
|
}:
|
||||||
case <-l.ctx.Done():
|
case <-scanCtx.Done():
|
||||||
return l.ctx.Err()
|
return scanCtx.Err()
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
@@ -321,8 +353,8 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
absolutePath: absoluteFilePath,
|
absolutePath: absoluteFilePath,
|
||||||
fileType: fileType,
|
fileType: fileType,
|
||||||
}:
|
}:
|
||||||
case <-l.ctx.Done():
|
case <-scanCtx.Done():
|
||||||
return l.ctx.Err()
|
return scanCtx.Err()
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
@@ -473,6 +505,10 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
|
|
||||||
for work := range workChan {
|
for work := range workChan {
|
||||||
g.Go(func() error {
|
g.Go(func() error {
|
||||||
|
if err := l.waitIfPaused(scanCtx); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
result, err := l.extractAudioMetadata(
|
result, err := l.extractAudioMetadata(
|
||||||
work, metrics,
|
work, metrics,
|
||||||
)
|
)
|
||||||
@@ -493,8 +529,8 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
|
|
||||||
select {
|
select {
|
||||||
case resultChan <- result:
|
case resultChan <- result:
|
||||||
case <-l.ctx.Done():
|
case <-scanCtx.Done():
|
||||||
return l.ctx.Err()
|
return scanCtx.Err()
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
@@ -546,74 +582,85 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
|
|
||||||
metrics.ThumbnailWallClock = time.Since(thumbStart)
|
metrics.ThumbnailWallClock = time.Since(thumbStart)
|
||||||
|
|
||||||
// --- Phase 5: orphan cleanup ---
|
// Skip orphan cleanup if the scan was cancelled — existingPaths
|
||||||
runtime.EventsEmit(l.ctx, events.LibraryScanProgress, ScanProgress{
|
// still contains unvisited files that would be incorrectly deleted.
|
||||||
Phase: "orphans", Total: totalFiles,
|
cancelled := scanCtx.Err() != nil
|
||||||
Processed: a + s + u, Added: a, Skipped: s, Updated: u,
|
|
||||||
})
|
|
||||||
|
|
||||||
orphanStart := time.Now()
|
|
||||||
|
|
||||||
var removed atomic.Int64
|
var removed atomic.Int64
|
||||||
|
|
||||||
existingPaths.Range(func(key, value any) bool {
|
if cancelled {
|
||||||
path := key.(string)
|
metrics.Cancelled = true
|
||||||
audioFile := value.(sqlcgen.AudioFile)
|
l.logger.Info("scan cancelled, skipping orphan cleanup")
|
||||||
|
} else {
|
||||||
|
// --- Phase 5: orphan cleanup ---
|
||||||
|
runtime.EventsEmit(l.ctx, events.LibraryScanProgress, ScanProgress{
|
||||||
|
Phase: "orphans", Total: totalFiles,
|
||||||
|
Processed: a + s + u, Added: a, Skipped: s, Updated: u,
|
||||||
|
})
|
||||||
|
|
||||||
l.logger.Debug(
|
orphanStart := time.Now()
|
||||||
"removing orphaned database entry",
|
|
||||||
"path", path, "id", audioFile.ID,
|
|
||||||
)
|
|
||||||
|
|
||||||
if err := l.db.Queries.DeleteAudioFile(
|
existingPaths.Range(func(key, value any) bool {
|
||||||
l.ctx, audioFile.ID,
|
path := key.(string)
|
||||||
); err != nil {
|
audioFile := value.(sqlcgen.AudioFile)
|
||||||
l.logger.Warn(
|
|
||||||
"failed to delete orphaned audio file",
|
l.logger.Debug(
|
||||||
"path", path,
|
"removing orphaned database entry",
|
||||||
"id", audioFile.ID,
|
"path", path, "id", audioFile.ID,
|
||||||
"err", err,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
metrics.addWarning(path, "orphan", err)
|
if err := l.db.Queries.DeleteAudioFile(
|
||||||
|
l.ctx, audioFile.ID,
|
||||||
|
); err != nil {
|
||||||
|
l.logger.Warn(
|
||||||
|
"failed to delete orphaned audio file",
|
||||||
|
"path", path,
|
||||||
|
"id", audioFile.ID,
|
||||||
|
"err", err,
|
||||||
|
)
|
||||||
|
|
||||||
|
metrics.addWarning(path, "orphan", err)
|
||||||
|
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
// Remove from FTS5 search index.
|
||||||
|
if err := l.db.DeleteSearchIndex(
|
||||||
|
audioFile.ID,
|
||||||
|
); err != nil {
|
||||||
|
l.logger.Warn(
|
||||||
|
"failed to delete FTS entry for orphan",
|
||||||
|
"id", audioFile.ID,
|
||||||
|
"err", err,
|
||||||
|
)
|
||||||
|
|
||||||
|
metrics.addWarning(path, "orphan", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
removed.Add(1)
|
||||||
|
|
||||||
return true
|
return true
|
||||||
}
|
})
|
||||||
|
|
||||||
// Remove from FTS5 search index.
|
metrics.OrphanCleanup = time.Since(orphanStart)
|
||||||
if err := l.db.DeleteSearchIndex(
|
}
|
||||||
audioFile.ID,
|
|
||||||
); err != nil {
|
// --- Phase 6: post-scan variant generation ---
|
||||||
|
if !cancelled {
|
||||||
|
variantStart := time.Now()
|
||||||
|
|
||||||
|
if err := l.generateMissingSizedVariants(); err != nil {
|
||||||
l.logger.Warn(
|
l.logger.Warn(
|
||||||
"failed to delete FTS entry for orphan",
|
"could not generate missing sized variants",
|
||||||
"id", audioFile.ID,
|
|
||||||
"err", err,
|
"err", err,
|
||||||
)
|
)
|
||||||
|
|
||||||
metrics.addWarning(path, "orphan", err)
|
metrics.addWarning("", "variant", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
removed.Add(1)
|
metrics.PostScanVariants = time.Since(variantStart)
|
||||||
|
|
||||||
return true
|
|
||||||
})
|
|
||||||
|
|
||||||
metrics.OrphanCleanup = time.Since(orphanStart)
|
|
||||||
|
|
||||||
// --- Phase 6: post-scan variant generation ---
|
|
||||||
variantStart := time.Now()
|
|
||||||
|
|
||||||
if err := l.generateMissingSizedVariants(); err != nil {
|
|
||||||
l.logger.Warn(
|
|
||||||
"could not generate missing sized variants",
|
|
||||||
"err", err,
|
|
||||||
)
|
|
||||||
|
|
||||||
metrics.addWarning("", "variant", err)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
metrics.PostScanVariants = time.Since(variantStart)
|
|
||||||
|
|
||||||
// --- Finalize ---
|
// --- Finalize ---
|
||||||
metrics.Added = added.Load()
|
metrics.Added = added.Load()
|
||||||
metrics.Updated = updated.Load()
|
metrics.Updated = updated.Load()
|
||||||
@@ -627,13 +674,18 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
"updated", metrics.Updated,
|
"updated", metrics.Updated,
|
||||||
"removed", metrics.Removed,
|
"removed", metrics.Removed,
|
||||||
"skipped", metrics.Skipped,
|
"skipped", metrics.Skipped,
|
||||||
|
"cancelled", cancelled,
|
||||||
"total", metrics.Total,
|
"total", metrics.Total,
|
||||||
"library", l.conf.DirectoryPath,
|
"library", l.conf.DirectoryPath,
|
||||||
)
|
)
|
||||||
|
|
||||||
runtime.EventsEmit(
|
if cancelled {
|
||||||
l.ctx, events.LibraryScanComplete, metrics,
|
runtime.EventsEmit(l.ctx, events.LibraryScanCancelled, metrics)
|
||||||
)
|
} else {
|
||||||
|
runtime.EventsEmit(
|
||||||
|
l.ctx, events.LibraryScanComplete, metrics,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
return metrics, scanErr
|
return metrics, scanErr
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,88 @@
|
|||||||
|
package library
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
|
||||||
|
"github.com/wailsapp/wails/v2/pkg/runtime"
|
||||||
|
|
||||||
|
"yellowjacket/backend/events"
|
||||||
|
)
|
||||||
|
|
||||||
|
// CancelScan cancels an in-progress scan. Returns immediately;
|
||||||
|
// scan goroutines stop at their next checkpoint.
|
||||||
|
func (l *Library) CancelScan() {
|
||||||
|
l.mu.Lock()
|
||||||
|
cancel := l.scanCancel
|
||||||
|
l.mu.Unlock()
|
||||||
|
|
||||||
|
if cancel != nil {
|
||||||
|
cancel()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// PauseScan pauses an in-progress scan. Workers block at their
|
||||||
|
// next pause checkpoint until ResumeScan is called.
|
||||||
|
func (l *Library) PauseScan() {
|
||||||
|
l.mu.Lock()
|
||||||
|
defer l.mu.Unlock()
|
||||||
|
|
||||||
|
if !l.scanActive || l.scanPaused {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
l.scanPaused = true
|
||||||
|
l.scanPauseCh = make(chan struct{})
|
||||||
|
|
||||||
|
runtime.EventsEmit(l.ctx, events.LibraryScanPaused)
|
||||||
|
}
|
||||||
|
|
||||||
|
// ResumeScan unblocks a paused scan.
|
||||||
|
func (l *Library) ResumeScan() {
|
||||||
|
l.mu.Lock()
|
||||||
|
defer l.mu.Unlock()
|
||||||
|
|
||||||
|
if !l.scanPaused {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
l.scanPaused = false
|
||||||
|
close(l.scanPauseCh) // unblocks all waiting workers
|
||||||
|
|
||||||
|
runtime.EventsEmit(l.ctx, events.LibraryScanResumed)
|
||||||
|
}
|
||||||
|
|
||||||
|
// IsScanActive returns whether a scan is currently running.
|
||||||
|
func (l *Library) IsScanActive() bool {
|
||||||
|
l.mu.Lock()
|
||||||
|
defer l.mu.Unlock()
|
||||||
|
|
||||||
|
return l.scanActive
|
||||||
|
}
|
||||||
|
|
||||||
|
// IsScanPaused returns whether the scan is currently paused.
|
||||||
|
func (l *Library) IsScanPaused() bool {
|
||||||
|
l.mu.Lock()
|
||||||
|
defer l.mu.Unlock()
|
||||||
|
|
||||||
|
return l.scanPaused
|
||||||
|
}
|
||||||
|
|
||||||
|
// waitIfPaused blocks the calling goroutine if the scan is paused.
|
||||||
|
// Returns ctx.Err() if the context is cancelled while waiting.
|
||||||
|
func (l *Library) waitIfPaused(ctx context.Context) error {
|
||||||
|
l.mu.Lock()
|
||||||
|
ch := l.scanPauseCh
|
||||||
|
paused := l.scanPaused
|
||||||
|
l.mu.Unlock()
|
||||||
|
|
||||||
|
if !paused || ch == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-ch: // closed = unpaused
|
||||||
|
return nil
|
||||||
|
case <-ctx.Done():
|
||||||
|
return ctx.Err()
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user