Files
bookhoard/TASKS-backend-progress-tracking.md
T
john-okeefe 88844670af Fix integration test bug: premature response body close
Fixed critical bug in TestScanProgress_BatchingWorks integration test where
response body was closed before JSON decoding, causing test failure.

Bug Location: Line 604 in scanner_integration_test.go example code

Problem:
  scanResp, err := client.Do(scanReq)
  require.NoError(s.T(), err)
  scanResp.Body.Close()  //  Closed here

  var scanResponse map[string]interface{}
  json.NewDecoder(scanResp.Body).Decode(&scanResponse)  //  Reads from closed body

Fix:
  scanResp, err := client.Do(scanReq)
  require.NoError(s.T(), err)

  var scanResponse map[string]interface{}
  json.NewDecoder(scanResp.Body).Decode(&scanResponse)
  scanResp.Body.Close()  //  Close AFTER decoding

This matches the pattern used in TestScanProgress_TracksStatistics and ensures
the response body is available for JSON decoding before being closed.

The implementation plan is now fully correct and ready for execution.
2026-02-25 10:47:14 -05:00

24 KiB

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:

{
  "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:

type JobResult struct {
    JobID    string
    Status   JobStatus
    Error    string
    Result   interface{}
    Progress float64
}

Add scan statistics:

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:

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:

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:

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:

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:

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:

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:

func (s *MediaScanner) GetStats() (int, int, int) {
    return s.totalFiles, s.newItems, s.errors
}

Update ScanFolders() to track stats:

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:

// 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:

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:

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:

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:

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

  • 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