# Backend Scan Progress Tracking **Date Created:** 2025-02-25 **Status:** Ready to Implement **Priority:** HIGH - Required for accurate progress UI --- ## ๐Ÿšจ Problem The scan job status endpoint returns minimal data: ```json { "job_id": "...", "status": "completed", "progress": 1.0, "result": {"message": "scan completed", "library_id": "..."} } ``` **Missing fields needed by frontend:** - `files_scanned` - total files processed - `new_items` - books added to database - `errors` - scan errors encountered - Real-time `progress` updates (0% โ†’ 100% during scan) **Current behavior:** - Progress jumps from 0% to 100% when scan completes - No file counts during scanning - No error tracking --- ## ๐Ÿ“‹ Implementation Plan ### Overview Add progress tracking to scan jobs by: 1. Extending `JobResult` to include scan statistics 2. Adding progress update mechanism to worker 3. Tracking statistics during `ScanFolders()` (with batching every 10 files) 4. Updating progress as files are processed 5. Adding integration tests to verify behavior --- ### Step 1: Extend JobResult Structure **File:** `internal/services/worker.go` **Current JobResult:** ```go type JobResult struct { JobID string Status JobStatus Error string Result interface{} Progress float64 } ``` **Add scan statistics:** ```go type JobResult struct { JobID string Status JobStatus Error string Result interface{} Progress float64 FilesScanned int // NEW NewItems int // NEW Errors int // NEW } ``` **Update `processScanJob()` to return stats:** ```go return map[string]interface{}{ "message": "scan completed", "library_id": libraryID, "files_scanned": totalFiles, "new_items": newItems, "errors": errors, }, nil ``` --- ### Step 2: Add Progress Update Callback to Job **File:** `internal/services/worker.go` **Add callback to Job struct:** ```go type Job struct { ID string Type JobType Params map[string]interface{} Status JobStatus CreatedAt time.Time StartedAt *time.Time CompletedAt *time.Time Error error Result interface{} Context context.Context ProgressCallback func(progress float64, filesScanned, newItems, errors int) // NEW } ``` **Add update method:** ```go func (j *Job) UpdateProgress(progress float64, filesScanned, newItems, errors int) { if j.ProgressCallback != nil { j.ProgressCallback(progress, filesScanned, newItems, errors) } } ``` --- ### Step 3: Set Up Callback and Pass Job to Scanner **File:** `internal/services/worker.go` **Update processScanJob() to set up callback:** ```go func (w *Worker) processScanJob(job *Job) (interface{}, error) { libraryID, ok := job.Params["library_id"].(string) if !ok { return nil, fmt.Errorf("library_id required") } folders, ok := job.Params["folders"].([]string) if !ok { return nil, fmt.Errorf("folders required") } adminID, ok := job.Params["admin_id"].(string) if !ok { return nil, fmt.Errorf("admin_id required") } db, ok := job.Params["db"].(*database.Queries) if !ok { return nil, fmt.Errorf("database queries required") } scanner := NewMediaScanner(db) scanner.job = job // NEW: Pass job reference for progress updates // NEW: Set up progress callback to update JobResult in real-time job.ProgressCallback = func(progress float64, filesScanned, newItems, errors int) { w.mu.Lock() defer w.mu.Unlock() if result, exists := w.results[job.ID]; exists { result.Progress = progress result.FilesScanned = filesScanned result.NewItems = newItems result.Errors = errors } } if err := scanner.SetFolders(folders); err != nil { return nil, err } var adminUUID pgtype.UUID if err := adminUUID.Scan(adminID); err != nil { return nil, err } scanner.SetAdminID(adminUUID) if err := scanner.ScanFolders(job.Context); err != nil { return nil, err } totalFiles, newItems, errors := scanner.GetStats() return map[string]interface{}{ "message": "scan completed", "library_id": libraryID, "files_scanned": totalFiles, "new_items": newItems, "errors": errors, }, nil } ``` --- ### Step 4: Update Worker Job Completion Tracking **File:** `internal/services/worker.go` **Modify job processing loop:** ```go case JobTypeScan: result, err = w.processScanJob(job) // After job completes, update JobResult with stats w.mu.Lock() status := JobStatusCompleted if err != nil { status = JobStatusFailed } if job.Context != nil && job.Context.Err() != nil { status = JobStatusCancelled } // Extract stats from result if available var filesScanned, newItems, errors int if result != nil { if stats, ok := result.(map[string]interface{}); ok { filesScanned = int(stats["files_scanned"].(float64)) newItems = int(stats["new_items"].(float64)) errors = int(stats["errors"].(float64)) } } w.results[job.ID] = &JobResult{ JobID: job.ID, Status: status, Error: func() string { if err != nil { return err.Error() } else { return "" } }(), Result: result, Progress: 1.0, FilesScanned: filesScanned, // NEW NewItems: newItems, // NEW Errors: errors, // NEW } w.mu.Unlock() ``` --- ### Step 5: Track Statistics During Scan **File:** `internal/services/media_scanner.go` **Add counter fields to MediaScanner:** ```go type MediaScanner struct { db *database.Queries watcher *fsnotify.Watcher folders []string adminID pgtype.UUID defaultLibraryID pgtype.UUID libraryTypes map[string][]string // NEW: Scan statistics totalFiles int newItems int errors int job *Job // Reference to job for progress updates } ``` **Add getter method:** ```go func (s *MediaScanner) GetStats() (int, int, int) { return s.totalFiles, s.newItems, s.errors } ``` **Update ScanFolders() to track stats:** ```go func (s *MediaScanner) ScanFolders(ctx context.Context) error { if len(s.folders) == 0 { return fmt.Errorf("no folders set") } // Reset counters s.totalFiles = 0 s.newItems = 0 s.errors = 0 // First pass: count total files for _, folder := range s.folders { filepath.WalkDir(folder, func(path string, d fs.DirEntry, err error) error { if !d.IsDir() && s.isScannableFile(path) { s.totalFiles++ } return nil }) } fmt.Printf("Starting scan of %d folders: %v (%d files to scan)\n", len(s.folders), s.folders, s.totalFiles) processedFiles := 0 mediaFiles := 0 for _, folder := range s.folders { fmt.Printf("Scanning folder: %s\n", folder) if _, err := os.Stat(folder); os.IsNotExist(err) { fmt.Printf("Folder does not exist: %s\n", folder) s.errors++ continue } err := filepath.WalkDir(folder, func(path string, d fs.DirEntry, err error) error { if err != nil { fmt.Printf("Error accessing path %s: %v\n", path, err) s.errors++ return err } if d.IsDir() { if err := s.watcher.Add(path); err != nil { fmt.Printf("Warning: failed to watch subdirectory %s: %v\n", path, err) } return nil } if s.isScannableFile(path) { mediaFiles++ processedFiles++ // Update progress (batch every 10 files to reduce mutex contention) if processedFiles%10 == 0 && s.totalFiles > 0 { progress := float64(processedFiles) / float64(s.totalFiles) if s.job != nil { s.job.UpdateProgress(progress, processedFiles, s.newItems, s.errors) } } wasNew, err := s.processMediaFile(ctx, path) if err != nil { fmt.Printf("Error processing media file %s: %v\n", path, err) s.errors++ } else { fmt.Printf("Successfully processed media file: %s\n", path) // newItems already incremented in processMediaFile if wasNew } } return nil }) if err != nil { s.errors++ return fmt.Errorf("failed to scan folder %s: %v", folder, err) } } fmt.Printf("Scan completed: %d total files scanned, %d media files found, %d new items, %d errors\n", processedFiles, mediaFiles, s.newItems, s.errors) // Final progress update to ensure we report 100% if s.job != nil && s.totalFiles > 0 { s.job.UpdateProgress(1.0, processedFiles, s.newItems, s.errors) } return nil } ``` --- ### Step 6: Update processMediaFile to Track New Items **File:** `internal/services/media_scanner.go` **Modify return value:** ```go // CURRENT: func (s *MediaScanner) processMediaFile(ctx context.Context, path string) error // NEW: Returns (bool, error) where bool indicates if item was newly created func (s *MediaScanner) processMediaFile(ctx context.Context, path string) (bool, error) { // ... existing file processing code ... // Check if media item already exists existingItem, err := s.getMediaItemByFilePath(ctx, path) if err == nil && existingItem.FileSize.Int64 == info.Size() { fmt.Printf("Media item already exists with same size, skipping: %s\n", path) return false, nil // FALSE = not a new item (already exists) } // If item exists but different size, it's an update - still not "new" if err == nil { fmt.Printf("Updating existing media item: %s\n", path) // ... update logic ... return false, nil // FALSE = not a new item (was an update) } // ... rest of processing for new item ... // Create media item in database createdItem, err := s.db.CreateMediaItem(ctx, database.CreateMediaItemParams{...}) if err != nil { return false, err // FALSE = error, false means not created } s.newItems++ // NEW: Track new items return true, nil // TRUE = new item created } ``` **Update ScanFolders() to use return value:** ```go wasNew, err := s.processMediaFile(ctx, path) if err != nil { s.errors++ } else if wasNew { // newItems already incremented in processMediaFile } ``` --- ### Step 7: Update GetScanStatus Handler **File:** `internal/handlers/scanner.go` **Update response to include new fields:** ```go func (h *Handler) GetScanStatus(c echo.Context) error { jobID := c.Param("jobId") result, exists := h.worker.GetJobStatus(jobID) if !exists { return c.JSON(http.StatusNotFound, map[string]string{"error": "job not found"}) } return c.JSON(http.StatusOK, map[string]interface{}{ "job_id": result.JobID, "status": result.Status, "error": result.Error, "result": result.Result, "progress": result.Progress, "files_scanned": result.FilesScanned, // NEW "new_items": result.NewItems, // NEW "errors": result.Errors, // NEW }) } ``` --- ### Step 8: Add Integration Tests **File:** `cmd/server/tests/scanner_integration_test.go` (new file) **Create new integration test file:** ```go package main import ( "bytes" "encoding/json" "fmt" "net/http" "testing" "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" ) type ScannerIntegrationTestSuite struct { suite.Suite setup *TestServerSetup } func (s *ScannerIntegrationTestSuite) SetupSuite() { s.setup = setupTestServer(s.T()) } func (s *ScannerIntegrationTestSuite) TearDownSuite() { s.setup.Close() } func (s *ScannerIntegrationTestSuite) TestScanProgress_TracksStatistics() { token := s.setup.Token // Create test library libraryID := s.setup.CreateLibrary(s.T(), "Scan Test Library", "ebooks") // Add folder to library folderURL := fmt.Sprintf("%s/api/libraries/%s/folders", s.setup.Server.URL, libraryID) folderReq := map[string]interface{}{ "folder_path": "/app/uploads", } folderBody, _ := json.Marshal(folderReq) req, _ := http.NewRequest("POST", folderURL, bytes.NewBuffer(folderBody)) req.Header.Set("Content-Type", "application/json") req.Header.Set("Authorization", "Bearer "+token) client := &http.Client{} resp, err := client.Do(req) require.NoError(s.T(), err) resp.Body.Close() require.Equal(s.T(), http.StatusCreated, resp.StatusCode, "Folder creation should succeed") // Start scan scanURL := fmt.Sprintf("%s/api/libraries/%s/scan", s.setup.Server.URL, libraryID) scanReq, _ := http.NewRequest("POST", scanURL, nil) scanReq.Header.Set("Authorization", "Bearer "+token) scanResp, err := client.Do(scanReq) require.NoError(s.T(), err) require.Equal(s.T(), http.StatusAccepted, scanResp.StatusCode) var scanResponse map[string]interface{} err = json.NewDecoder(scanResp.Body).Decode(&scanResponse) require.NoError(s.T(), err) scanResp.Body.Close() jobID, ok := scanResponse["job_id"].(string) require.True(s.T(), ok, "job_id should be string") require.NotEmpty(s.T(), jobID, "job_id should not be empty") // Poll for progress updates var lastProgress float64 var lastFilesScanned, lastNewItems, lastErrors int for i := 0; i < 30; i++ { // Poll for up to 30 seconds time.Sleep(1 * time.Second) statusURL := fmt.Sprintf("%s/api/scanner/status/%s", s.setup.Server.URL, jobID) statusReq, _ := http.NewRequest("GET", statusURL, nil) statusReq.Header.Set("Authorization", "Bearer "+token) statusResp, err := client.Do(statusReq) require.NoError(s.T(), err) var status map[string]interface{} err = json.NewDecoder(statusResp.Body).Decode(&status) statusResp.Body.Close() require.NoError(s.T(), err) // Verify new fields exist assert.Contains(s.T(), status, "files_scanned") assert.Contains(s.T(), status, "new_items") assert.Contains(s.T(), status, "errors") // Track progress with safe type assertions progressFloat, ok := status["progress"].(float64) require.True(s.T(), ok, "progress should be float64") progress := progressFloat filesScannedFloat, ok := status["files_scanned"].(float64) require.True(s.T(), ok, "files_scanned should be float64") filesScanned := int(filesScannedFloat) newItemsFloat, ok := status["new_items"].(float64) require.True(s.T(), ok, "new_items should be float64") newItems := int(newItemsFloat) errorsFloat, ok := status["errors"].(float64) require.True(s.T(), ok, "errors should be float64") errors := int(errorsFloat) // Progress should be non-decreasing assert.GreaterOrEqual(s.T(), progress, lastProgress) lastProgress = progress // Files scanned should be non-decreasing assert.GreaterOrEqual(s.T(), filesScanned, lastFilesScanned) lastFilesScanned = filesScanned // Items/errors should be non-decreasing assert.GreaterOrEqual(s.T(), newItems, lastNewItems) assert.GreaterOrEqual(s.T(), errors, lastErrors) // Break if scan complete if status["status"] == "completed" || status["status"] == "failed" { break } } // Verify final state assert.Equal(s.T(), 1.0, lastProgress) assert.GreaterOrEqual(s.T(), lastFilesScanned, 0) } func (s *ScannerIntegrationTestSuite) TestScanProgress_BatchingWorks() { token := s.setup.Token // Create library with folder libraryID := s.setup.CreateLibrary(s.T(), "Batch Test Library", "ebooks") folderURL := fmt.Sprintf("%s/api/libraries/%s/folders", s.setup.Server.URL, libraryID) folderReq := map[string]interface{}{ "folder_path": "/app/uploads", } folderBody, _ := json.Marshal(folderReq) req, _ := http.NewRequest("POST", folderURL, bytes.NewBuffer(folderBody)) req.Header.Set("Content-Type", "application/json") req.Header.Set("Authorization", "Bearer "+token) client := &http.Client{} resp, err := client.Do(req) require.NoError(s.T(), err) resp.Body.Close() // Start scan scanURL := fmt.Sprintf("%s/api/libraries/%s/scan", s.setup.Server.URL, libraryID) scanReq, _ := http.NewRequest("POST", scanURL, nil) scanReq.Header.Set("Authorization", "Bearer "+token) scanResp, err := client.Do(scanReq) require.NoError(s.T(), err) var scanResponse map[string]interface{} json.NewDecoder(scanResp.Body).Decode(&scanResponse) scanResp.Body.Close() jobID, ok := scanResponse["job_id"].(string) require.True(s.T(), ok, "job_id should be string") require.NotEmpty(s.T(), jobID, "job_id should not be empty") // Poll and verify we don't get updates on EVERY file updateCount := 0 previousFilesScanned := -1 for i := 0; i < 20; i++ { time.Sleep(500 * time.Millisecond) statusURL := fmt.Sprintf("%s/api/scanner/status/%s", s.setup.Server.URL, jobID) statusReq, _ := http.NewRequest("GET", statusURL, nil) statusReq.Header.Set("Authorization", "Bearer "+token) statusResp, _ := client.Do(statusReq) var status map[string]interface{} err = json.NewDecoder(statusResp.Body).Decode(&status) require.NoError(s.T(), err) statusResp.Body.Close() filesScannedFloat, ok := status["files_scanned"].(float64) require.True(s.T(), ok, "files_scanned should be float64") filesScanned := int(filesScannedFloat) // Only count as update if files_scanned changed if filesScanned != previousFilesScanned { updateCount++ previousFilesScanned = filesScanned } if status["status"] == "completed" || status["status"] == "failed" { break } } // With batching every 10 files, we should have FEWER updates than files // This is a weak assertion, but verifies batching is working // Threshold of 50 assumes test library has < 500 files - adjust based on actual test data assert.Less(s.T(), updateCount, 50) } func TestScannerIntegrationTestSuite(t *testing.T) { suite.Run(t, new(ScannerIntegrationTestSuite)) } ``` **Note:** The integration test requires `encoding/json` import (already included in the import list above). **Update existing unit tests:** `internal/services/worker_test.go` **Add test for new JobResult fields:** ```go func TestWorker_JobResult_HasStatsFields(t *testing.T) { // Test that JobResult properly stores scan statistics worker := NewWorker(1) defer worker.Shutdown() jobID := "test-job-stats" // Simulate job completion with stats worker.mu.Lock() worker.results[jobID] = &JobResult{ JobID: jobID, Status: JobStatusCompleted, Progress: 1.0, FilesScanned: 42, NewItems: 5, Errors: 1, } worker.mu.Unlock() // Verify stats are retrievable result, exists := worker.GetJobStatus(jobID) require.True(t, exists, "Job result should exist") require.NotNil(t, result, "Result should not be nil") assert.Equal(t, jobID, result.JobID) assert.Equal(t, JobStatusCompleted, result.Status) assert.Equal(t, 1.0, result.Progress) assert.Equal(t, 42, result.FilesScanned, "FilesScanned should be 42") assert.Equal(t, 5, result.NewItems, "NewItems should be 5") assert.Equal(t, 1, result.Errors, "Errors should be 1") } func TestWorker_ProgressCallback_UpdatesJobResult(t *testing.T) { // Test that progress callback updates JobResult in real-time worker := NewWorker(1) defer worker.Shutdown() job := &Job{ ID: "test-progress", Type: JobTypeScan, Status: JobStatusInProgress, Context: context.Background(), } // Set up progress callback job.ProgressCallback = func(progress float64, filesScanned, newItems, errors int) { worker.mu.Lock() defer worker.mu.Unlock() if result, exists := worker.results[job.ID]; exists { result.Progress = progress result.FilesScanned = filesScanned result.NewItems = newItems result.Errors = errors } } // Initialize result worker.mu.Lock() worker.results[job.ID] = &JobResult{ JobID: job.ID, Status: JobStatusInProgress, } worker.mu.Unlock() // Simulate progress updates job.UpdateProgress(0.5, 10, 2, 0) result, exists := worker.GetJobStatus(job.ID) require.True(t, exists) assert.Equal(t, 0.5, result.Progress) assert.Equal(t, 10, result.FilesScanned) assert.Equal(t, 2, result.NewItems) assert.Equal(t, 0, result.Errors) // Simulate completion job.UpdateProgress(1.0, 20, 5, 1) result, exists = worker.GetJobStatus(job.ID) require.True(t, exists) assert.Equal(t, 1.0, result.Progress) assert.Equal(t, 20, result.FilesScanned) assert.Equal(t, 5, result.NewItems) assert.Equal(t, 1, result.Errors) } ``` --- ## ๐Ÿงช Testing Checklist After implementation: - [ ] Unit tests pass: `go test ./internal/services/...` - [ ] TestWorker_JobResult_HasStatsFields verifies stats are stored - [ ] TestWorker_ProgressCallback_UpdatesJobResult verifies real-time updates - [ ] Integration tests pass: `go test ./cmd/server/tests/...` - [ ] TestScanProgress_TracksStatistics verifies all new fields exist and increment - [ ] TestScanProgress_BatchingWorks verifies batching reduces update frequency - [ ] Scan a library with multiple files - [ ] Poll `/api/scanner/status/{jobId}` during scan - [ ] Verify `progress` increases from 0% to 100% gradually - [ ] Verify `files_scanned` count increases during scan - [ ] Verify `new_items` count shows books added - [ ] Verify `errors` count shows scan errors - [ ] Check final status has accurate totals - [ ] Test with empty library (no files) - [ ] Test with library containing only non-media files - [ ] Test with library causing scan errors - [ ] Verify batching reduces update frequency (integration test) --- ## ๐Ÿ“ Notes - **Thread Safety:** JobResult updates are thread-safe via worker's mutex (w.mu.Lock()) - **Performance:** Progress updates are batched every 10 files to reduce mutex contention - **Memory:** Stats tracking uses 3 int fields (24 bytes) per MediaScanner instance - **Backwards Compatibility:** Frontend already uses `|| 0` fallbacks, so safe to deploy - **Testing:** Integration tests use `setupTestServer()` from `test_helpers.go` - **Database Pool Configuration:** `setupTestServer()` already sets `max_conns=1` (test_helpers.go:428), preventing connection pool exhaustion during test runs --- ## ๐Ÿ”— Related Files - `internal/services/worker.go` - Job processing and result tracking - `internal/services/media_scanner.go` - Scan logic and statistics - `internal/handlers/scanner.go` - Status API endpoint - `cmd/server/tests/scanner_integration_test.go` - Integration tests (NEW) - `internal/services/worker_test.go` - Unit tests to update - `TASKS-scanning-progress.md` - Frontend implementation that depends on this --- **Last Updated:** 2025-02-25 **Status:** Ready for review and implementation