feat(11-01): per-library scan pipeline with queue coordinator
Task 1: Schema, events, and progress types - Add library_id to CreateAudioFile SQL INSERT and regenerate sqlc code - Add LibraryScanQueued and LibraryScanQueueDrained event constants - Regenerate TypeScript events via genevents - Add LibraryID, LibraryName, QueuedCount to ScanProgress - Add LibraryID, LibraryName to ScanMetrics - Add libraryID field to importResult for threading through pipeline Task 2: Scan queue coordinator and per-library scanning - Create scan_queue.go with ScanLibrary(id), ScanAllLibraries() - Add CancelCurrentScan(), CancelAllScans() for queue-aware cancellation - FIFO scan queue with silent dedup (same library already scanning or queued) - Refactor Scan() -> scanInternal(libraryID, libraryName, libraryPath) - Replace GetAllAudioFiles with GetAudioFilesByLibrary for per-library loading - Thread libraryID through DB writer to set CreateAudioFileParams.LibraryID - drainQueue auto-starts next queued library or emits LibraryScanQueueDrained - Pause freezes current scan AND queue - Add GetScanQueueLength() and QueuedLibraryNames() for UI - Mark CancelScan() and Scan() as deprecated
This commit is contained in:
@@ -1,5 +1,5 @@
|
|||||||
-- name: CreateAudioFile :one
|
-- name: CreateAudioFile :one
|
||||||
INSERT INTO audio_files (file_path, length_milliseconds, file_type_id, recording_id, sample_rate, bit_depth, channels, bitrate, file_size, basename) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
INSERT INTO audio_files (file_path, length_milliseconds, file_type_id, recording_id, sample_rate, bit_depth, channels, bitrate, file_size, basename, library_id) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
RETURNING *;
|
RETURNING *;
|
||||||
|
|
||||||
-- name: GetAudioFile :one
|
-- name: GetAudioFile :one
|
||||||
|
|||||||
@@ -34,7 +34,7 @@ func (q *Queries) CountAudioFilesByLibrary(ctx context.Context, libraryID int64)
|
|||||||
}
|
}
|
||||||
|
|
||||||
const createAudioFile = `-- name: CreateAudioFile :one
|
const createAudioFile = `-- name: CreateAudioFile :one
|
||||||
INSERT INTO audio_files (file_path, length_milliseconds, file_type_id, recording_id, sample_rate, bit_depth, channels, bitrate, file_size, basename) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
INSERT INTO audio_files (file_path, length_milliseconds, file_type_id, recording_id, sample_rate, bit_depth, channels, bitrate, file_size, basename, library_id) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||||
RETURNING id, file_path, length_milliseconds, file_type_id, recording_id, sample_rate, bit_depth, channels, bitrate, file_size, basename, library_id
|
RETURNING id, file_path, length_milliseconds, file_type_id, recording_id, sample_rate, bit_depth, channels, bitrate, file_size, basename, library_id
|
||||||
`
|
`
|
||||||
|
|
||||||
@@ -49,6 +49,7 @@ type CreateAudioFileParams struct {
|
|||||||
Bitrate int64
|
Bitrate int64
|
||||||
FileSize int64
|
FileSize int64
|
||||||
Basename string
|
Basename string
|
||||||
|
LibraryID int64
|
||||||
}
|
}
|
||||||
|
|
||||||
func (q *Queries) CreateAudioFile(ctx context.Context, arg CreateAudioFileParams) (AudioFile, error) {
|
func (q *Queries) CreateAudioFile(ctx context.Context, arg CreateAudioFileParams) (AudioFile, error) {
|
||||||
@@ -63,6 +64,7 @@ func (q *Queries) CreateAudioFile(ctx context.Context, arg CreateAudioFileParams
|
|||||||
arg.Bitrate,
|
arg.Bitrate,
|
||||||
arg.FileSize,
|
arg.FileSize,
|
||||||
arg.Basename,
|
arg.Basename,
|
||||||
|
arg.LibraryID,
|
||||||
)
|
)
|
||||||
var i AudioFile
|
var i AudioFile
|
||||||
err := row.Scan(
|
err := row.Scan(
|
||||||
|
|||||||
@@ -54,3 +54,9 @@ const (
|
|||||||
LibraryScanPaused = "LibraryScanPaused"
|
LibraryScanPaused = "LibraryScanPaused"
|
||||||
LibraryScanResumed = "LibraryScanResumed"
|
LibraryScanResumed = "LibraryScanResumed"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// Scan queue events.
|
||||||
|
const (
|
||||||
|
LibraryScanQueued = "LibraryScanQueued"
|
||||||
|
LibraryScanQueueDrained = "LibraryScanQueueDrained"
|
||||||
|
)
|
||||||
|
|||||||
+114
-51
@@ -90,6 +90,11 @@ type Library struct {
|
|||||||
scanCancel context.CancelFunc
|
scanCancel context.CancelFunc
|
||||||
scanPaused bool
|
scanPaused bool
|
||||||
scanPauseCh chan struct{}
|
scanPauseCh chan struct{}
|
||||||
|
|
||||||
|
// Scan queue fields — protected by mu.
|
||||||
|
scanQueue []scanQueueEntry
|
||||||
|
currentScanLibraryID int64
|
||||||
|
currentScanLibraryName string
|
||||||
}
|
}
|
||||||
|
|
||||||
// SetRescanHooks provides optional hooks for cross-cutting
|
// SetRescanHooks provides optional hooks for cross-cutting
|
||||||
@@ -174,12 +179,40 @@ func (l *Library) registerEventHandlers() {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
// Scan syncs the library by adding new files and removing deleted ones.
|
// Scan syncs the library using the legacy DirectoryPath config.
|
||||||
// Files that exist but have incomplete metadata (recording_id = 0)
|
// Retained for backward compatibility with handleConfigUpdate.
|
||||||
// will be updated. The returned ScanMetrics contains timing and
|
//
|
||||||
// count data for every phase of the scan.
|
// Deprecated: Use ScanLibrary(id) for per-library scanning.
|
||||||
func (l *Library) Scan() (*ScanMetrics, error) {
|
func (l *Library) Scan() (*ScanMetrics, error) {
|
||||||
|
if len(l.conf.DirectoryPath) == 0 {
|
||||||
|
return newScanMetrics(), errLibraryDirNotConfigured
|
||||||
|
}
|
||||||
|
|
||||||
|
lib, err := l.db.Queries.GetLibraryByPath(
|
||||||
|
l.ctx, string(l.conf.DirectoryPath),
|
||||||
|
)
|
||||||
|
if err != nil {
|
||||||
|
return newScanMetrics(), fmt.Errorf(
|
||||||
|
"could not resolve library for path %s: %w",
|
||||||
|
l.conf.DirectoryPath, err,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
return l.scanInternal(lib.ID, lib.Name, lib.Path), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// scanInternal performs the full scan pipeline for a single library.
|
||||||
|
// It is called from the scan queue coordinator (startScan) or the
|
||||||
|
// legacy Scan() wrapper. The caller is responsible for goroutine
|
||||||
|
// management; this method blocks until the scan completes.
|
||||||
|
func (l *Library) scanInternal(
|
||||||
|
libraryID int64,
|
||||||
|
libraryName string,
|
||||||
|
libraryPath string,
|
||||||
|
) *ScanMetrics {
|
||||||
metrics := newScanMetrics()
|
metrics := newScanMetrics()
|
||||||
|
metrics.LibraryID = libraryID
|
||||||
|
metrics.LibraryName = libraryName
|
||||||
scanStart := time.Now()
|
scanStart := time.Now()
|
||||||
|
|
||||||
scanCtx, scanCancel := context.WithCancel(l.ctx)
|
scanCtx, scanCancel := context.WithCancel(l.ctx)
|
||||||
@@ -195,7 +228,6 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
defer func() {
|
defer func() {
|
||||||
l.mu.Lock()
|
l.mu.Lock()
|
||||||
l.scanCancel = nil
|
l.scanCancel = nil
|
||||||
l.scanActive = false
|
|
||||||
// If still paused, unpause so no dangling channel.
|
// If still paused, unpause so no dangling channel.
|
||||||
if l.scanPaused {
|
if l.scanPaused {
|
||||||
l.scanPaused = false
|
l.scanPaused = false
|
||||||
@@ -208,28 +240,54 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
l.mu.Unlock()
|
l.mu.Unlock()
|
||||||
}()
|
}()
|
||||||
|
|
||||||
if len(l.conf.DirectoryPath) == 0 {
|
|
||||||
return metrics, errLibraryDirNotConfigured
|
|
||||||
}
|
|
||||||
|
|
||||||
workerCount := resolveScanWorkerCount(
|
workerCount := resolveScanWorkerCount(
|
||||||
l.conf.ScanConcurrency,
|
ScanConcurrencyAuto,
|
||||||
string(l.conf.DirectoryPath),
|
libraryPath,
|
||||||
)
|
)
|
||||||
|
|
||||||
l.logger.Info(
|
l.logger.Info(
|
||||||
"beginning library scan",
|
"beginning library scan",
|
||||||
|
"libraryID", libraryID,
|
||||||
|
"libraryName", libraryName,
|
||||||
|
"libraryPath", libraryPath,
|
||||||
"workers", workerCount,
|
"workers", workerCount,
|
||||||
"concurrencyMode", l.conf.ScanConcurrency,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
runtime.EventsEmit(l.ctx, events.LibraryScanStarted)
|
// Helper to build a ScanProgress with library identification.
|
||||||
|
queuedCount := func() int {
|
||||||
|
l.mu.Lock()
|
||||||
|
defer l.mu.Unlock()
|
||||||
|
|
||||||
basePath := string(l.conf.DirectoryPath)
|
return len(l.scanQueue)
|
||||||
|
}
|
||||||
|
|
||||||
|
mkProgress := func(
|
||||||
|
phase string,
|
||||||
|
total, processed, a, s, u int64,
|
||||||
|
) ScanProgress {
|
||||||
|
return ScanProgress{
|
||||||
|
Phase: phase,
|
||||||
|
Total: total,
|
||||||
|
Processed: processed,
|
||||||
|
Added: a,
|
||||||
|
Skipped: s,
|
||||||
|
Updated: u,
|
||||||
|
LibraryID: libraryID,
|
||||||
|
LibraryName: libraryName,
|
||||||
|
QueuedCount: queuedCount(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
runtime.EventsEmit(l.ctx, events.LibraryScanStarted, map[string]any{
|
||||||
|
"libraryId": libraryID,
|
||||||
|
"libraryName": libraryName,
|
||||||
|
})
|
||||||
|
|
||||||
|
basePath := libraryPath
|
||||||
|
|
||||||
// --- Pre-walk: count audio files for progress reporting ---
|
// --- Pre-walk: count audio files for progress reporting ---
|
||||||
runtime.EventsEmit(l.ctx, events.LibraryScanProgress,
|
runtime.EventsEmit(l.ctx, events.LibraryScanProgress,
|
||||||
ScanProgress{Phase: "counting"},
|
mkProgress("counting", 0, 0, 0, 0, 0),
|
||||||
)
|
)
|
||||||
|
|
||||||
totalFiles := countAudioFiles(basePath)
|
totalFiles := countAudioFiles(basePath)
|
||||||
@@ -239,14 +297,20 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
"total", totalFiles,
|
"total", totalFiles,
|
||||||
)
|
)
|
||||||
|
|
||||||
// --- Phase 1: load existing files from DB ---
|
// --- Phase 1: load existing files from DB (per-library) ---
|
||||||
loadStart := time.Now()
|
loadStart := time.Now()
|
||||||
|
|
||||||
existingFiles, err := l.db.Queries.GetAllAudioFiles(l.ctx)
|
existingFiles, err := l.db.Queries.GetAudioFilesByLibrary(
|
||||||
|
l.ctx, libraryID,
|
||||||
|
)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return metrics, fmt.Errorf(
|
l.logger.Error(
|
||||||
"could not load existing audio files: %w", err,
|
"could not load existing audio files",
|
||||||
|
"libraryID", libraryID,
|
||||||
|
"err", err,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
return metrics
|
||||||
}
|
}
|
||||||
|
|
||||||
existingPaths := &sync.Map{}
|
existingPaths := &sync.Map{}
|
||||||
@@ -259,7 +323,8 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
l.logger.Debug(
|
l.logger.Debug(
|
||||||
"loaded existing files from database",
|
"loaded existing files from database",
|
||||||
"count", len(existingFiles),
|
"count", len(existingFiles),
|
||||||
"library-directory", l.conf.DirectoryPath,
|
"libraryID", libraryID,
|
||||||
|
"libraryPath", libraryPath,
|
||||||
)
|
)
|
||||||
|
|
||||||
workChan := make(chan scanWork, 100)
|
workChan := make(chan scanWork, 100)
|
||||||
@@ -424,14 +489,10 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
runtime.EventsEmit(
|
runtime.EventsEmit(
|
||||||
l.ctx,
|
l.ctx,
|
||||||
events.LibraryScanProgress,
|
events.LibraryScanProgress,
|
||||||
ScanProgress{
|
mkProgress(
|
||||||
Phase: "scanning",
|
"scanning", totalFiles,
|
||||||
Total: totalFiles,
|
a+s+u, a, s, u,
|
||||||
Processed: a + s + u,
|
),
|
||||||
Added: a,
|
|
||||||
Skipped: s,
|
|
||||||
Updated: u,
|
|
||||||
},
|
|
||||||
)
|
)
|
||||||
case <-stopProgress:
|
case <-stopProgress:
|
||||||
return
|
return
|
||||||
@@ -477,6 +538,9 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for result := range resultChan {
|
for result := range resultChan {
|
||||||
|
// Thread library ID into each result for saveAudioFile.
|
||||||
|
result.libraryID = libraryID
|
||||||
|
|
||||||
if !dbStarted {
|
if !dbStarted {
|
||||||
dbStartVal = time.Now()
|
dbStartVal = time.Now()
|
||||||
dbStarted = true
|
dbStarted = true
|
||||||
@@ -553,14 +617,7 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
u := updated.Load()
|
u := updated.Load()
|
||||||
|
|
||||||
runtime.EventsEmit(l.ctx, events.LibraryScanProgress,
|
runtime.EventsEmit(l.ctx, events.LibraryScanProgress,
|
||||||
ScanProgress{
|
mkProgress("scanning", totalFiles, a+s+u, a, s, u),
|
||||||
Phase: "scanning",
|
|
||||||
Total: totalFiles,
|
|
||||||
Processed: a + s + u,
|
|
||||||
Added: a,
|
|
||||||
Skipped: s,
|
|
||||||
Updated: u,
|
|
||||||
},
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Close thumbnail channel and wait for all thumbnail workers
|
// Close thumbnail channel and wait for all thumbnail workers
|
||||||
@@ -568,14 +625,9 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
// point so it is safe to close.
|
// point so it is safe to close.
|
||||||
thumbStart := time.Now()
|
thumbStart := time.Now()
|
||||||
|
|
||||||
runtime.EventsEmit(l.ctx, events.LibraryScanProgress, ScanProgress{
|
runtime.EventsEmit(l.ctx, events.LibraryScanProgress,
|
||||||
Phase: "thumbnails",
|
mkProgress("thumbnails", totalFiles, a+s+u, a, s, u),
|
||||||
Total: totalFiles,
|
)
|
||||||
Processed: a + s + u,
|
|
||||||
Added: a,
|
|
||||||
Skipped: s,
|
|
||||||
Updated: u,
|
|
||||||
})
|
|
||||||
|
|
||||||
close(thumbChan)
|
close(thumbChan)
|
||||||
thumbWg.Wait()
|
thumbWg.Wait()
|
||||||
@@ -594,10 +646,9 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
l.logger.Info("scan cancelled, skipping orphan cleanup")
|
l.logger.Info("scan cancelled, skipping orphan cleanup")
|
||||||
} else {
|
} else {
|
||||||
// --- Phase 5: orphan cleanup ---
|
// --- Phase 5: orphan cleanup ---
|
||||||
runtime.EventsEmit(l.ctx, events.LibraryScanProgress, ScanProgress{
|
runtime.EventsEmit(l.ctx, events.LibraryScanProgress,
|
||||||
Phase: "orphans", Total: totalFiles,
|
mkProgress("orphans", totalFiles, a+s+u, a, s, u),
|
||||||
Processed: a + s + u, Added: a, Skipped: s, Updated: u,
|
)
|
||||||
})
|
|
||||||
|
|
||||||
orphanStart := time.Now()
|
orphanStart := time.Now()
|
||||||
|
|
||||||
@@ -669,26 +720,36 @@ func (l *Library) Scan() (*ScanMetrics, error) {
|
|||||||
metrics.Removed = removed.Load()
|
metrics.Removed = removed.Load()
|
||||||
metrics.Total = time.Since(scanStart)
|
metrics.Total = time.Since(scanStart)
|
||||||
|
|
||||||
|
if scanErr != nil {
|
||||||
|
l.logger.Warn(
|
||||||
|
"scan completed with errors",
|
||||||
|
"err", scanErr,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
l.logger.Info(
|
l.logger.Info(
|
||||||
"library scan complete",
|
"library scan complete",
|
||||||
|
"libraryID", libraryID,
|
||||||
|
"libraryName", libraryName,
|
||||||
"added", metrics.Added,
|
"added", metrics.Added,
|
||||||
"updated", metrics.Updated,
|
"updated", metrics.Updated,
|
||||||
"removed", metrics.Removed,
|
"removed", metrics.Removed,
|
||||||
"skipped", metrics.Skipped,
|
"skipped", metrics.Skipped,
|
||||||
"cancelled", cancelled,
|
"cancelled", cancelled,
|
||||||
"total", metrics.Total,
|
"total", metrics.Total,
|
||||||
"library", l.conf.DirectoryPath,
|
|
||||||
)
|
)
|
||||||
|
|
||||||
if cancelled {
|
if cancelled {
|
||||||
runtime.EventsEmit(l.ctx, events.LibraryScanCancelled, metrics)
|
runtime.EventsEmit(
|
||||||
|
l.ctx, events.LibraryScanCancelled, metrics,
|
||||||
|
)
|
||||||
} else {
|
} else {
|
||||||
runtime.EventsEmit(
|
runtime.EventsEmit(
|
||||||
l.ctx, events.LibraryScanComplete, metrics,
|
l.ctx, events.LibraryScanComplete, metrics,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
return metrics, scanErr
|
return metrics
|
||||||
}
|
}
|
||||||
|
|
||||||
// progressInterval controls how often scan progress events are
|
// progressInterval controls how often scan progress events are
|
||||||
@@ -765,6 +826,7 @@ type importResult struct {
|
|||||||
audioProps *metadata.AudioProperties
|
audioProps *metadata.AudioProperties
|
||||||
existingFileID int64 // non-zero if this is an update
|
existingFileID int64 // non-zero if this is an update
|
||||||
needsUpdate bool
|
needsUpdate bool
|
||||||
|
libraryID int64 // library this file belongs to
|
||||||
}
|
}
|
||||||
|
|
||||||
// extractAudioMetadata reads and extracts metadata from an audio file.
|
// extractAudioMetadata reads and extracts metadata from an audio file.
|
||||||
@@ -939,6 +1001,7 @@ func (l *Library) saveAudioFile(
|
|||||||
Bitrate: int64(props.Bitrate),
|
Bitrate: int64(props.Bitrate),
|
||||||
FileSize: props.FileSize,
|
FileSize: props.FileSize,
|
||||||
Basename: basename,
|
Basename: basename,
|
||||||
|
LibraryID: result.libraryID,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf(
|
return fmt.Errorf(
|
||||||
|
|||||||
@@ -53,6 +53,10 @@ type ScanMetrics struct {
|
|||||||
// Cancelled is true when the scan was stopped via CancelScan.
|
// Cancelled is true when the scan was stopped via CancelScan.
|
||||||
Cancelled bool `json:"cancelled"`
|
Cancelled bool `json:"cancelled"`
|
||||||
|
|
||||||
|
// Library identification.
|
||||||
|
LibraryID int64 `json:"libraryId"` // library that was scanned
|
||||||
|
LibraryName string `json:"libraryName"` // display name of scanned library
|
||||||
|
|
||||||
// Non-fatal issues encountered during scanning.
|
// Non-fatal issues encountered during scanning.
|
||||||
Warnings []ScanWarning `json:"warnings"`
|
Warnings []ScanWarning `json:"warnings"`
|
||||||
}
|
}
|
||||||
@@ -60,12 +64,15 @@ type ScanMetrics struct {
|
|||||||
// ScanProgress is the payload emitted periodically during a scan to
|
// ScanProgress is the payload emitted periodically during a scan to
|
||||||
// report live progress to the frontend.
|
// report live progress to the frontend.
|
||||||
type ScanProgress struct {
|
type ScanProgress struct {
|
||||||
Phase string `json:"phase"` // "counting", "scanning", "orphans", "thumbnails"
|
Phase string `json:"phase"` // "counting", "scanning", "orphans", "thumbnails"
|
||||||
Total int64 `json:"total"` // total audio files from pre-walk count
|
Total int64 `json:"total"` // total audio files from pre-walk count
|
||||||
Processed int64 `json:"processed"` // added + skipped + updated so far
|
Processed int64 `json:"processed"` // added + skipped + updated so far
|
||||||
Added int64 `json:"added"`
|
Added int64 `json:"added"`
|
||||||
Skipped int64 `json:"skipped"`
|
Skipped int64 `json:"skipped"`
|
||||||
Updated int64 `json:"updated"`
|
Updated int64 `json:"updated"`
|
||||||
|
LibraryID int64 `json:"libraryId"` // library being scanned
|
||||||
|
LibraryName string `json:"libraryName"` // display name of library being scanned
|
||||||
|
QueuedCount int `json:"queuedCount"` // number of libraries still queued after this one
|
||||||
}
|
}
|
||||||
|
|
||||||
// ScanWarning represents a non-fatal issue encountered during scanning.
|
// ScanWarning represents a non-fatal issue encountered during scanning.
|
||||||
|
|||||||
@@ -10,6 +10,9 @@ import (
|
|||||||
|
|
||||||
// CancelScan cancels an in-progress scan. Returns immediately;
|
// CancelScan cancels an in-progress scan. Returns immediately;
|
||||||
// scan goroutines stop at their next checkpoint.
|
// scan goroutines stop at their next checkpoint.
|
||||||
|
//
|
||||||
|
// Deprecated: Use CancelCurrentScan or CancelAllScans for
|
||||||
|
// queue-aware cancellation.
|
||||||
func (l *Library) CancelScan() {
|
func (l *Library) CancelScan() {
|
||||||
l.mu.Lock()
|
l.mu.Lock()
|
||||||
cancel := l.scanCancel
|
cancel := l.scanCancel
|
||||||
|
|||||||
@@ -0,0 +1,173 @@
|
|||||||
|
package library
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
|
||||||
|
"github.com/wailsapp/wails/v2/pkg/runtime"
|
||||||
|
|
||||||
|
"yellowjacket/backend/events"
|
||||||
|
)
|
||||||
|
|
||||||
|
// scanQueueEntry holds the metadata needed to scan a single library.
|
||||||
|
type scanQueueEntry struct {
|
||||||
|
libraryID int64
|
||||||
|
libraryName string
|
||||||
|
libraryPath string
|
||||||
|
}
|
||||||
|
|
||||||
|
// ScanLibrary queues a scan for the library with the given database ID.
|
||||||
|
// If no scan is active the library is scanned immediately; otherwise it
|
||||||
|
// is appended to the queue. Duplicate requests (same library already
|
||||||
|
// scanning or already queued) are silently ignored.
|
||||||
|
func (l *Library) ScanLibrary(id int64) error {
|
||||||
|
lib, err := l.db.Queries.GetLibrary(l.ctx, id)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("could not get library %d: %w", id, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
l.mu.Lock()
|
||||||
|
defer l.mu.Unlock()
|
||||||
|
|
||||||
|
// Silent dedup: already scanning this library.
|
||||||
|
if l.currentScanLibraryID == id {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Silent dedup: already queued.
|
||||||
|
for _, entry := range l.scanQueue {
|
||||||
|
if entry.libraryID == id {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
entry := scanQueueEntry{
|
||||||
|
libraryID: lib.ID,
|
||||||
|
libraryName: lib.Name,
|
||||||
|
libraryPath: lib.Path,
|
||||||
|
}
|
||||||
|
|
||||||
|
if !l.scanActive {
|
||||||
|
l.scanActive = true
|
||||||
|
l.currentScanLibraryID = entry.libraryID
|
||||||
|
l.currentScanLibraryName = entry.libraryName
|
||||||
|
|
||||||
|
go l.startScan(entry)
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// A scan is already running — queue this library.
|
||||||
|
l.scanQueue = append(l.scanQueue, entry)
|
||||||
|
|
||||||
|
runtime.EventsEmit(l.ctx, events.LibraryScanQueued, map[string]any{
|
||||||
|
"libraryId": lib.ID,
|
||||||
|
"libraryName": lib.Name,
|
||||||
|
"queueLength": len(l.scanQueue),
|
||||||
|
})
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ScanAllLibraries queries all libraries from the database and queues
|
||||||
|
// each one for scanning. Existing dedup logic ensures no duplicates.
|
||||||
|
func (l *Library) ScanAllLibraries() error {
|
||||||
|
libs, err := l.db.Queries.GetAllLibraries(l.ctx)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("could not get all libraries: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, lib := range libs {
|
||||||
|
if err := l.ScanLibrary(lib.ID); err != nil {
|
||||||
|
l.logger.Warn(
|
||||||
|
"could not queue library for scan",
|
||||||
|
"libraryID", lib.ID,
|
||||||
|
"libraryName", lib.Name,
|
||||||
|
"err", err,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// CancelCurrentScan cancels only the currently scanning library.
|
||||||
|
// The next queued library (if any) starts automatically when the
|
||||||
|
// current scan's goroutine completes.
|
||||||
|
func (l *Library) CancelCurrentScan() {
|
||||||
|
l.mu.Lock()
|
||||||
|
cancel := l.scanCancel
|
||||||
|
l.mu.Unlock()
|
||||||
|
|
||||||
|
if cancel != nil {
|
||||||
|
cancel()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// CancelAllScans cancels the current scan and clears the entire
|
||||||
|
// queue so no further libraries are scanned.
|
||||||
|
func (l *Library) CancelAllScans() {
|
||||||
|
l.mu.Lock()
|
||||||
|
l.scanQueue = nil
|
||||||
|
cancel := l.scanCancel
|
||||||
|
l.mu.Unlock()
|
||||||
|
|
||||||
|
if cancel != nil {
|
||||||
|
cancel()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// GetScanQueueLength returns the number of libraries waiting in the
|
||||||
|
// scan queue (excludes the currently scanning library).
|
||||||
|
func (l *Library) GetScanQueueLength() int {
|
||||||
|
l.mu.Lock()
|
||||||
|
defer l.mu.Unlock()
|
||||||
|
|
||||||
|
return len(l.scanQueue)
|
||||||
|
}
|
||||||
|
|
||||||
|
// QueuedLibraryNames returns the display names of libraries waiting
|
||||||
|
// in the scan queue, in FIFO order.
|
||||||
|
func (l *Library) QueuedLibraryNames() []string {
|
||||||
|
l.mu.Lock()
|
||||||
|
defer l.mu.Unlock()
|
||||||
|
|
||||||
|
names := make([]string, len(l.scanQueue))
|
||||||
|
for i, entry := range l.scanQueue {
|
||||||
|
names[i] = entry.libraryName
|
||||||
|
}
|
||||||
|
|
||||||
|
return names
|
||||||
|
}
|
||||||
|
|
||||||
|
// startScan runs the scan for a single library entry and then drains
|
||||||
|
// the queue. It is always called in a new goroutine.
|
||||||
|
func (l *Library) startScan(entry scanQueueEntry) {
|
||||||
|
l.scanInternal(entry.libraryID, entry.libraryName, entry.libraryPath)
|
||||||
|
l.drainQueue()
|
||||||
|
}
|
||||||
|
|
||||||
|
// drainQueue is called after each scan completes. If the queue is
|
||||||
|
// non-empty the next entry is popped and scanned; otherwise the
|
||||||
|
// scan pipeline is marked idle.
|
||||||
|
func (l *Library) drainQueue() {
|
||||||
|
l.mu.Lock()
|
||||||
|
|
||||||
|
if len(l.scanQueue) > 0 {
|
||||||
|
next := l.scanQueue[0]
|
||||||
|
l.scanQueue = l.scanQueue[1:]
|
||||||
|
l.currentScanLibraryID = next.libraryID
|
||||||
|
l.currentScanLibraryName = next.libraryName
|
||||||
|
l.mu.Unlock()
|
||||||
|
|
||||||
|
go l.startScan(next)
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
l.currentScanLibraryID = 0
|
||||||
|
l.currentScanLibraryName = ""
|
||||||
|
l.scanActive = false
|
||||||
|
l.mu.Unlock()
|
||||||
|
|
||||||
|
runtime.EventsEmit(l.ctx, events.LibraryScanQueueDrained)
|
||||||
|
}
|
||||||
@@ -38,6 +38,10 @@ export const Events = {
|
|||||||
LibraryScanCancelled: "LibraryScanCancelled",
|
LibraryScanCancelled: "LibraryScanCancelled",
|
||||||
LibraryScanPaused: "LibraryScanPaused",
|
LibraryScanPaused: "LibraryScanPaused",
|
||||||
LibraryScanResumed: "LibraryScanResumed",
|
LibraryScanResumed: "LibraryScanResumed",
|
||||||
|
|
||||||
|
// Scan queue events
|
||||||
|
LibraryScanQueued: "LibraryScanQueued",
|
||||||
|
LibraryScanQueueDrained: "LibraryScanQueueDrained",
|
||||||
} as const;
|
} as const;
|
||||||
|
|
||||||
export type EventName = (typeof Events)[keyof typeof Events];
|
export type EventName = (typeof Events)[keyof typeof Events];
|
||||||
|
|||||||
+12
@@ -3,6 +3,10 @@
|
|||||||
import {library} from '../models';
|
import {library} from '../models';
|
||||||
import {context} from '../models';
|
import {context} from '../models';
|
||||||
|
|
||||||
|
export function CancelAllScans():Promise<void>;
|
||||||
|
|
||||||
|
export function CancelCurrentScan():Promise<void>;
|
||||||
|
|
||||||
export function CancelScan():Promise<void>;
|
export function CancelScan():Promise<void>;
|
||||||
|
|
||||||
export function FullRescan():Promise<library.ScanMetrics>;
|
export function FullRescan():Promise<library.ScanMetrics>;
|
||||||
@@ -19,6 +23,8 @@ export function GetAllGenresWithCounts():Promise<Array<library.GenreWithCount>>;
|
|||||||
|
|
||||||
export function GetAllTracks():Promise<Array<library.Track>>;
|
export function GetAllTracks():Promise<Array<library.Track>>;
|
||||||
|
|
||||||
|
export function GetScanQueueLength():Promise<number>;
|
||||||
|
|
||||||
export function GetTracksByGenre(arg1:string):Promise<Array<library.Track>>;
|
export function GetTracksByGenre(arg1:string):Promise<Array<library.Track>>;
|
||||||
|
|
||||||
export function IsScanActive():Promise<boolean>;
|
export function IsScanActive():Promise<boolean>;
|
||||||
@@ -27,10 +33,16 @@ export function IsScanPaused():Promise<boolean>;
|
|||||||
|
|
||||||
export function PauseScan():Promise<void>;
|
export function PauseScan():Promise<void>;
|
||||||
|
|
||||||
|
export function QueuedLibraryNames():Promise<Array<string>>;
|
||||||
|
|
||||||
export function ResumeScan():Promise<void>;
|
export function ResumeScan():Promise<void>;
|
||||||
|
|
||||||
export function Scan():Promise<library.ScanMetrics>;
|
export function Scan():Promise<library.ScanMetrics>;
|
||||||
|
|
||||||
|
export function ScanAllLibraries():Promise<void>;
|
||||||
|
|
||||||
|
export function ScanLibrary(arg1:number):Promise<void>;
|
||||||
|
|
||||||
export function SearchTracks(arg1:string):Promise<Array<library.Track>>;
|
export function SearchTracks(arg1:string):Promise<Array<library.Track>>;
|
||||||
|
|
||||||
export function SetContext(arg1:context.Context):Promise<void>;
|
export function SetContext(arg1:context.Context):Promise<void>;
|
||||||
|
|||||||
@@ -2,6 +2,14 @@
|
|||||||
// Cynhyrchwyd y ffeil hon yn awtomatig. PEIDIWCH Â MODIWL
|
// Cynhyrchwyd y ffeil hon yn awtomatig. PEIDIWCH Â MODIWL
|
||||||
// This file is automatically generated. DO NOT EDIT
|
// This file is automatically generated. DO NOT EDIT
|
||||||
|
|
||||||
|
export function CancelAllScans() {
|
||||||
|
return window['go']['library']['Library']['CancelAllScans']();
|
||||||
|
}
|
||||||
|
|
||||||
|
export function CancelCurrentScan() {
|
||||||
|
return window['go']['library']['Library']['CancelCurrentScan']();
|
||||||
|
}
|
||||||
|
|
||||||
export function CancelScan() {
|
export function CancelScan() {
|
||||||
return window['go']['library']['Library']['CancelScan']();
|
return window['go']['library']['Library']['CancelScan']();
|
||||||
}
|
}
|
||||||
@@ -34,6 +42,10 @@ export function GetAllTracks() {
|
|||||||
return window['go']['library']['Library']['GetAllTracks']();
|
return window['go']['library']['Library']['GetAllTracks']();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export function GetScanQueueLength() {
|
||||||
|
return window['go']['library']['Library']['GetScanQueueLength']();
|
||||||
|
}
|
||||||
|
|
||||||
export function GetTracksByGenre(arg1) {
|
export function GetTracksByGenre(arg1) {
|
||||||
return window['go']['library']['Library']['GetTracksByGenre'](arg1);
|
return window['go']['library']['Library']['GetTracksByGenre'](arg1);
|
||||||
}
|
}
|
||||||
@@ -50,6 +62,10 @@ export function PauseScan() {
|
|||||||
return window['go']['library']['Library']['PauseScan']();
|
return window['go']['library']['Library']['PauseScan']();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export function QueuedLibraryNames() {
|
||||||
|
return window['go']['library']['Library']['QueuedLibraryNames']();
|
||||||
|
}
|
||||||
|
|
||||||
export function ResumeScan() {
|
export function ResumeScan() {
|
||||||
return window['go']['library']['Library']['ResumeScan']();
|
return window['go']['library']['Library']['ResumeScan']();
|
||||||
}
|
}
|
||||||
@@ -58,6 +74,14 @@ export function Scan() {
|
|||||||
return window['go']['library']['Library']['Scan']();
|
return window['go']['library']['Library']['Scan']();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export function ScanAllLibraries() {
|
||||||
|
return window['go']['library']['Library']['ScanAllLibraries']();
|
||||||
|
}
|
||||||
|
|
||||||
|
export function ScanLibrary(arg1) {
|
||||||
|
return window['go']['library']['Library']['ScanLibrary'](arg1);
|
||||||
|
}
|
||||||
|
|
||||||
export function SearchTracks(arg1) {
|
export function SearchTracks(arg1) {
|
||||||
return window['go']['library']['Library']['SearchTracks'](arg1);
|
return window['go']['library']['Library']['SearchTracks'](arg1);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -109,6 +109,8 @@ export namespace library {
|
|||||||
skipped: number;
|
skipped: number;
|
||||||
removed: number;
|
removed: number;
|
||||||
cancelled: boolean;
|
cancelled: boolean;
|
||||||
|
libraryId: number;
|
||||||
|
libraryName: string;
|
||||||
warnings: ScanWarning[];
|
warnings: ScanWarning[];
|
||||||
|
|
||||||
static createFrom(source: any = {}) {
|
static createFrom(source: any = {}) {
|
||||||
@@ -143,6 +145,8 @@ export namespace library {
|
|||||||
this.skipped = source["skipped"];
|
this.skipped = source["skipped"];
|
||||||
this.removed = source["removed"];
|
this.removed = source["removed"];
|
||||||
this.cancelled = source["cancelled"];
|
this.cancelled = source["cancelled"];
|
||||||
|
this.libraryId = source["libraryId"];
|
||||||
|
this.libraryName = source["libraryName"];
|
||||||
this.warnings = this.convertValues(source["warnings"], ScanWarning);
|
this.warnings = this.convertValues(source["warnings"], ScanWarning);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user