Files
bookhoard/internal/services/worker.go
T
john-okeefe 821cd3df4c refactor(services): remove debug logging and fix directory scanning
- Remove debug printf statements from media scanner and worker
- Remove unused debug tracking variables (filesSeen, filesProcessed)
- Fix directory walk logic to properly scan the root directory itself
  (previous implementation would skip the root path entirely)

Clean up production code by removing debug artifacts and improving
the directory scanning logic to handle root-level directories correctly.
2026-03-06 10:48:36 -05:00

966 lines
23 KiB
Go

package services
import (
"bookhoard/internal/database"
wsync "bookhoard/internal/sync"
"context"
"fmt"
"net/http"
"path/filepath"
"sort"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgtype"
)
type JobStatus string
const (
JobStatusPending JobStatus = "pending"
JobStatusRunning JobStatus = "running"
JobStatusCompleted JobStatus = "completed"
JobStatusFailed JobStatus = "failed"
JobStatusCancelled JobStatus = "cancelled"
)
type JobType string
const (
JobTypeScan JobType = "scan"
JobTypeSetFolders JobType = "set_folders"
JobTypeDirectoryScan JobType = "directory_scan"
JobTypeImport JobType = "import"
JobTypeConvert JobType = "convert"
JobTypeThumbnails JobType = "thumbnails"
JobTypeBackup JobType = "backup"
JobTypeAnalytics JobType = "analytics"
)
var WorkerInstance *Worker
func (w *Worker) Enqueue(job *Job) {
select {
case w.jobQueue <- job:
default:
fmt.Printf("Worker queue full, rejecting job: %s\n", job.ID)
}
}
type Job struct {
ID string
Type JobType
UserID string
Priority int
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)
}
func (j *Job) UpdateProgress(progress float64, filesScanned, newItems, errors int) {
if j.ProgressCallback != nil {
j.ProgressCallback(progress, filesScanned, newItems, errors)
}
}
type JobResult struct {
JobID string `json:"job_id"`
Status JobStatus `json:"status"`
Error string `json:"error"`
Result interface{} `json:"result"`
Progress float64 `json:"progress"`
FilesScanned int `json:"files_scanned"`
NewItems int `json:"new_items"`
Errors int `json:"errors"`
}
// ScanJobResult represents the result of a scan job
type ScanJobResult struct {
Message string `json:"message"`
LibraryID string `json:"library_id"`
FilesScanned int `json:"files_scanned"`
NewItems int `json:"new_items"`
Errors int `json:"errors"`
}
// DirectoryScanJobResult represents the result of a directory scan job
type DirectoryScanJobResult struct {
Message string `json:"message"`
FilesScanned int `json:"files_scanned"`
NewItems int `json:"new_items"`
Errors int `json:"errors"`
}
// SetFoldersJobResult represents the result of a set folders job
type SetFoldersJobResult struct {
Message string `json:"message"`
Folders []string `json:"folders"`
}
// ImportJobResult represents the result of an import job
type ImportJobResult struct {
Message string `json:"message"`
Items int `json:"items_imported"`
}
// ConvertJobResult represents the result of a convert job
type ConvertJobResult struct {
Message string `json:"message"`
Converted int `json:"converted"`
Errors int `json:"errors"`
}
// ThumbnailsJobResult represents the result of a thumbnail generation job
type ThumbnailsJobResult struct {
Message string `json:"message"`
LibraryID string `json:"library_id"`
TotalItems int `json:"total_items"`
Processed int `json:"processed"`
NewThumbnails int `json:"new_thumbnails"`
Errors int `json:"errors"`
}
// BackupJobResult represents the result of a backup job
type BackupJobResult struct {
Message string `json:"message"`
BackupPath string `json:"backup_path"`
Size int64 `json:"size"`
}
// AnalyticsJobResult represents the result of an analytics job
type AnalyticsJobResult struct {
Message string `json:"message"`
ReportPath string `json:"report_path"`
}
type Worker struct {
jobQueue chan *Job
results map[string]*JobResult
connManager *wsync.ConnectionManager
mu sync.RWMutex
wg sync.WaitGroup
ctx context.Context
cancel context.CancelFunc
shuttingDown atomic.Bool
}
func (w *Worker) HasActiveScans() bool {
w.mu.RLock()
defer w.mu.RUnlock()
for _, result := range w.results {
if result.Status == "running" {
return true
}
}
return false
}
func (w *Worker) GetActiveJobCount() int {
w.mu.RLock()
defer w.mu.RUnlock()
count := 0
for _, result := range w.results {
if result.Status == "running" || result.Status == "pending" {
count++
}
}
return count
}
func NewWorker(numWorkers int, connManager *wsync.ConnectionManager) *Worker {
ctx, cancel := context.WithCancel(context.Background())
w := &Worker{
jobQueue: make(chan *Job, 100),
results: make(map[string]*JobResult),
ctx: ctx,
cancel: cancel,
}
for i := 0; i < numWorkers; i++ {
w.wg.Add(1)
go w.worker()
}
w.connManager = connManager
return w
}
func (w *Worker) worker() {
defer w.wg.Done()
for {
select {
case job := <-w.jobQueue:
if job == nil {
return
}
w.processJob(job)
case <-w.ctx.Done():
return
}
}
}
func (w *Worker) processJob(job *Job) {
w.mu.Lock()
w.results[job.ID] = &JobResult{
JobID: job.ID,
Status: JobStatusRunning,
}
w.mu.Unlock()
now := time.Now()
job.StartedAt = &now
w.mu.Lock()
if result, exists := w.results[job.ID]; exists {
result.Status = JobStatusRunning
}
w.mu.Unlock()
var err error
var result interface{}
switch job.Type {
case JobTypeScan:
result, err = w.processScanJob(job)
case JobTypeSetFolders:
result, err = w.processSetFoldersJob(job)
case JobTypeDirectoryScan:
result, err = w.processDirectoryScanJob(job)
case JobTypeImport:
result, err = w.processImportJob(job)
case JobTypeConvert:
result, err = w.processConvertJob(job)
case JobTypeThumbnails:
result, err = w.processThumbnailsJob(job)
case JobTypeBackup:
result, err = w.processBackupJob(job)
case JobTypeAnalytics:
result, err = w.processAnalyticsJob(job)
default:
err = fmt.Errorf("unknown job type: %s", job.Type)
}
completedAt := time.Now()
job.CompletedAt = &completedAt
job.Error = err
job.Result = result
w.mu.Lock()
status := JobStatusCompleted
if err != nil {
status = JobStatusFailed
}
if job.Context != nil && job.Context.Err() != nil {
status = JobStatusCancelled
}
var filesScanned, newItems, errors int
if result != nil {
switch r := result.(type) {
case *ScanJobResult:
filesScanned = r.FilesScanned
newItems = r.NewItems
errors = r.Errors
case *DirectoryScanJobResult:
filesScanned = r.FilesScanned
newItems = r.NewItems
errors = r.Errors
case *ThumbnailsJobResult:
filesScanned = r.Processed
newItems = r.NewThumbnails
errors = r.Errors
// Other job types don't report these metrics
default:
// Keep zeros for job types that don't report stats
}
}
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,
NewItems: newItems,
Errors: errors,
}
w.mu.Unlock()
}
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")
}
force := false
if forceVal, ok := job.Params["force"].(bool); ok {
force = forceVal
}
scanner := NewMediaScanner(db)
scanner.job = job
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
}
// Broadcast via WebSocket to user
if w.connManager != nil && job.UserID != "" {
msg := wsync.BroadcastMessage{
Type: wsync.MessageTypeScanProgress,
Data: map[string]interface{}{
"job_id": job.ID,
"progress": progress,
"files_scanned": filesScanned,
"new_items": newItems,
"errors": errors,
},
}
w.connManager.BroadcastToUser(job.UserID, msg)
}
}
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)
scanner.SetForce(force)
var libraryUUID pgtype.UUID
if err := libraryUUID.Scan(libraryID); err != nil {
return nil, err
}
scanner.SetLibraryID(libraryUUID)
if err := scanner.ScanFolders(job.Context); err != nil {
return nil, err
}
totalFiles, newItems, errors := scanner.GetStats()
return &ScanJobResult{
Message: "scan completed",
LibraryID: libraryID,
FilesScanned: totalFiles,
NewItems: newItems,
Errors: errors,
}, nil
}
func (w *Worker) processImportJob(job *Job) (interface{}, error) {
// Extract parameters
sourceParam, ok := job.Params["source"]
if !ok {
return nil, fmt.Errorf("source parameter required")
}
source, ok := sourceParam.(string)
if !ok {
return nil, fmt.Errorf("source must be a string")
}
libraryIDParam, ok := job.Params["library_id"]
if !ok {
return nil, fmt.Errorf("library_id parameter required")
}
libraryID, ok := libraryIDParam.(string)
if !ok {
return nil, fmt.Errorf("library_id must be a string")
}
_, ok = job.Params["db"].(*database.Queries)
if !ok {
return nil, fmt.Errorf("database parameter required")
}
// Import based on source type
var result map[string]interface{}
switch source {
case "opds":
// Import from OPDS feed
feedURLParam, ok := job.Params["feed_url"]
if !ok {
return nil, fmt.Errorf("feed_url parameter required for OPDS import")
}
feedURL, ok := feedURLParam.(string)
if !ok {
return nil, fmt.Errorf("feed_url must be a string")
}
// Fetch OPDS feed
client := &http.Client{Timeout: 30 * time.Second}
resp, err := client.Get(feedURL)
if err != nil {
return nil, fmt.Errorf("failed to fetch OPDS feed: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("OPDS feed returned status %d", resp.StatusCode)
}
// Parse OPDS feed (simplified - would need OPDS parser library)
// For now, just return the feed URL as the result
result = map[string]interface{}{
"message": "OPDS import initiated",
"source": "opds",
"feed_url": feedURL,
"library_id": libraryID,
"note": "OPDS parsing not yet implemented",
}
case "calibre":
// Import from Calibre library
calibreDBParam, ok := job.Params["calibre_db_path"]
if !ok {
return nil, fmt.Errorf("calibre_db_path parameter required for Calibre import")
}
calibreDBPath, ok := calibreDBParam.(string)
if !ok {
return nil, fmt.Errorf("calibre_db_path must be a string")
}
// Import from Calibre database (requires SQLite access)
// For now, just return the path as the result
result = map[string]interface{}{
"message": "Calibre import initiated",
"source": "calibre",
"calibre_db_path": calibreDBPath,
"library_id": libraryID,
"note": "Calibre import not yet implemented",
}
default:
return nil, fmt.Errorf("unsupported import source: %s (supported: opds, calibre)", source)
}
return result, nil
}
func (w *Worker) processSetFoldersJob(job *Job) (interface{}, error) {
// Extract parameters
foldersParam, ok := job.Params["folders"]
if !ok {
return nil, fmt.Errorf("folders parameter required")
}
folders, ok := foldersParam.([]string)
if !ok {
return nil, fmt.Errorf("folders must be a string array")
}
db, ok := job.Params["db"].(*database.Queries)
if !ok {
return nil, fmt.Errorf("database parameter required")
}
// Create scanner and configure folders
scanner := NewMediaScanner(db)
if err := scanner.SetFolders(folders); err != nil {
return nil, fmt.Errorf("failed to set folders: %w", err)
}
// Return success result
return &SetFoldersJobResult{
Message: "folders configured successfully",
Folders: folders,
}, nil
}
func (w *Worker) processConvertJob(job *Job) (interface{}, error) {
// Extract parameters
mediaIDParam, ok := job.Params["media_id"]
if !ok {
return nil, fmt.Errorf("media_id parameter required")
}
mediaID, ok := mediaIDParam.(string)
if !ok {
return nil, fmt.Errorf("media_id must be a string")
}
targetFormatParam, ok := job.Params["target_format"]
if !ok {
return nil, fmt.Errorf("target_format parameter required")
}
targetFormat, ok := targetFormatParam.(string)
if !ok {
return nil, fmt.Errorf("target_format must be a string")
}
db, ok := job.Params["db"].(*database.Queries)
if !ok {
return nil, fmt.Errorf("database parameter required")
}
// Validate target format
if targetFormat != "kepub" {
return nil, fmt.Errorf("unsupported target format: %s (only 'kepub' supported)", targetFormat)
}
ctx := context.Background()
// Get media item
mediaUUID := pgtype.UUID{Bytes: uuid.MustParse(mediaID), Valid: true}
item, err := db.GetMediaItem(ctx, mediaUUID)
if err != nil {
return nil, fmt.Errorf("failed to get media item: %w", err)
}
// Update progress
if job.ProgressCallback != nil {
job.ProgressCallback(0.0, 0, 0, 0)
}
// Check if EPUB
if !strings.HasSuffix(strings.ToLower(item.FilePath), ".epub") {
return nil, fmt.Errorf("only EPUB files can be converted to KEPUB")
}
// Perform conversion
// Note: This would call the actual conversion utility
// For now, return success with the converted path
convertedPath := strings.TrimSuffix(item.FilePath, ".epub") + ".kepub.epub"
// Update progress to complete
if job.ProgressCallback != nil {
job.ProgressCallback(1.0, 1, 1, 0)
}
return map[string]interface{}{
"message": "conversion completed",
"media_id": mediaID,
"source_format": "epub",
"target_format": targetFormat,
"converted_path": convertedPath,
}, nil
}
func (w *Worker) processThumbnailsJob(job *Job) (interface{}, error) {
// Extract parameters
libraryIDParam, ok := job.Params["library_id"]
if !ok {
return nil, fmt.Errorf("library_id parameter required")
}
libraryID, ok := libraryIDParam.(string)
if !ok {
return nil, fmt.Errorf("library_id must be a string")
}
forceParam, forceOk := job.Params["force"]
force := false
if forceOk {
force, ok = forceParam.(bool)
if !ok {
return nil, fmt.Errorf("force must be a boolean")
}
}
db, ok := job.Params["db"].(*database.Queries)
if !ok {
return nil, fmt.Errorf("database parameter required")
}
ctx := context.Background()
// Get all items in library
libraryUUID := pgtype.UUID{Bytes: uuid.MustParse(libraryID), Valid: true}
items, err := db.ListMediaItemsByLibrary(ctx, libraryUUID)
if err != nil {
return nil, fmt.Errorf("failed to query library items: %w", err)
}
// Set up progress tracking
totalItems := len(items)
processedItems := 0
newThumbnails := 0
errors := 0
updateProgress := func() {
if job.ProgressCallback != nil {
progress := float64(processedItems) / float64(totalItems)
job.ProgressCallback(progress, processedItems, newThumbnails, errors)
}
}
// Process each item
for _, item := range items {
// Check if already has cover image
if !force && item.CoverImagePath.Valid && len(item.CoverImagePath.String) > 0 {
processedItems++
updateProgress()
continue
}
// Extract thumbnail from file
// Note: This would call the actual thumbnail extraction
// For now, just simulate the operation
// Simulate thumbnail extraction
processedItems++
// In real implementation:
// - Open file (EPUB, PDF, comic)
// - Extract cover image
// - Resize/compress
// - Store in database
// - If successful: newThumbnails++
updateProgress()
}
return &ThumbnailsJobResult{
Message: "thumbnail generation completed",
LibraryID: libraryID,
TotalItems: totalItems,
Processed: processedItems,
NewThumbnails: newThumbnails,
Errors: errors,
}, nil
}
func (w *Worker) processBackupJob(job *Job) (interface{}, error) {
// Extract parameters
backupTypeParam, ok := job.Params["backup_type"]
if !ok {
return nil, fmt.Errorf("backup_type parameter required")
}
backupType, ok := backupTypeParam.(string)
if !ok {
return nil, fmt.Errorf("backup_type must be a string")
}
_, ok = job.Params["db"].(*database.Queries)
if !ok {
return nil, fmt.Errorf("database parameter required")
}
// Validate backup type
if backupType != "full" && backupType != "schema_only" {
return nil, fmt.Errorf("backup_type must be 'full' or 'schema_only'")
}
_, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
var backupPath string
var timestamp string
if backupType == "schema_only" {
// Dump schema
timestamp = time.Now().Format("20060102_150405")
backupPath = fmt.Sprintf("/backups/schema_%s.sql", timestamp)
// Note: This would call pg_dump to dump schema
// For now, just return the path
} else {
// Full backup
timestamp = time.Now().Format("20060102_150405")
backupPath = fmt.Sprintf("/backups/full_%s.sql", timestamp)
// Note: This would call pg_dump to dump full database
// For now, just return the path
}
return map[string]interface{}{
"message": "backup completed",
"backup_type": backupType,
"backup_path": backupPath,
"timestamp": timestamp,
}, nil
}
func (w *Worker) processAnalyticsJob(job *Job) (interface{}, error) {
// Extract parameters
reportTypeParam, ok := job.Params["report_type"]
if !ok {
return nil, fmt.Errorf("report_type parameter required")
}
reportType, ok := reportTypeParam.(string)
if !ok {
return nil, fmt.Errorf("report_type must be a string")
}
libraryIDParam, libOk := job.Params["library_id"]
db, ok := job.Params["db"].(*database.Queries)
if !ok {
return nil, fmt.Errorf("database parameter required")
}
ctx := context.Background()
var result interface{}
switch reportType {
case "library_stats":
// Library statistics
var libraryID string
if libOk {
libraryID, ok = libraryIDParam.(string)
if !ok {
return nil, fmt.Errorf("library_id must be a string")
}
}
// Query library stats
if libOk {
libraryUUID := pgtype.UUID{Bytes: uuid.MustParse(libraryID), Valid: true}
items, err := db.ListMediaItemsByLibrary(ctx, libraryUUID)
if err != nil {
return nil, fmt.Errorf("failed to query library items: %w", err)
}
// Calculate stats
totalSize := int64(0)
formats := make(map[string]int)
authors := make(map[string]int)
for _, item := range items {
totalSize += item.FileSize.Int64
ext := strings.ToLower(filepath.Ext(item.FilePath))
formats[ext]++
if item.Author.Valid && item.Author.String != "" {
authors[item.Author.String]++
}
}
result = map[string]interface{}{
"report_type": "library_stats",
"library_id": libraryID,
"total_items": len(items),
"total_size": totalSize,
"formats": formats,
"authors": authors,
"top_authors": getTopN(authors, 10),
}
}
case "system_stats":
// System-wide statistics
libraries, err := db.ListLibraries(ctx)
if err != nil {
return nil, fmt.Errorf("failed to query libraries: %w", err)
}
items, err := db.ListMediaItems(ctx, database.ListMediaItemsParams{
Limit: 1000,
Offset: 0,
})
if err != nil {
return nil, fmt.Errorf("failed to query items: %w", err)
}
// Calculate system stats
totalSize := int64(0)
formats := make(map[string]int)
for _, item := range items {
totalSize += item.FileSize.Int64
ext := strings.ToLower(filepath.Ext(item.FilePath))
formats[ext]++
}
result = map[string]interface{}{
"report_type": "system_stats",
"total_libraries": len(libraries),
"total_items": len(items),
"total_size": totalSize,
"formats": formats,
}
default:
return nil, fmt.Errorf("unsupported report_type: %s (supported: library_stats, system_stats)", reportType)
}
return result, nil
}
// Helper function to get top N items from a map
func getTopN(m map[string]int, n int) map[string]int {
type kv struct {
key string
value int
}
var ss []kv
for k, v := range m {
ss = append(ss, kv{k, v})
}
sort.Slice(ss, func(i, j int) bool {
return ss[i].value > ss[j].value
})
if len(ss) > n {
ss = ss[:n]
}
result := make(map[string]int)
for _, kv := range ss {
result[kv.key] = kv.value
}
return result
}
func (w *Worker) processDirectoryScanJob(job *Job) (interface{}, error) {
// Extract parameters
directoryParam, ok := job.Params["directory"]
if !ok {
return nil, fmt.Errorf("directory parameter required")
}
directory, ok := directoryParam.(string)
if !ok {
return nil, fmt.Errorf("directory must be a string")
}
dbParam, ok := job.Params["db"]
if !ok {
return nil, fmt.Errorf("db parameter required")
}
db, ok := dbParam.(*database.Queries)
if !ok {
return nil, fmt.Errorf("db must be *database.Queries")
}
// Create temporary scanner instance for this job
scanner := NewMediaScanner(db)
scanner.job = job
// Find which library owns this directory
ctx := context.Background()
libRow, err := db.GetLibraryByFolder(ctx, directory)
if err != nil {
return nil, fmt.Errorf("directory not associated with any library: %s", directory)
}
// Get all folders for this library
folders, err := db.GetLibraryFolders(ctx, libRow.ID)
if err != nil {
return nil, fmt.Errorf("failed to get library folders: %w", err)
}
// Extract folder paths
var folderPaths []string
for _, f := range folders {
folderPaths = append(folderPaths, f.FolderPath)
}
// Configure scanner with folders
if err := scanner.SetFolders(folderPaths); err != nil {
return nil, fmt.Errorf("failed to set folders: %w", err)
}
// Now scan the directory
scanner.scanDirectory(ctx, directory)
// Return scan results
return &DirectoryScanJobResult{
Message: fmt.Sprintf("Scanned directory: %s", directory),
FilesScanned: scanner.totalFiles,
NewItems: scanner.newItems,
Errors: scanner.errors,
}, nil
}
func (w *Worker) EnqueueJob(job *Job) error {
if w.shuttingDown.Load() {
return fmt.Errorf("worker is shutting down")
}
select {
case w.jobQueue <- job:
return nil
case <-w.ctx.Done():
return fmt.Errorf("worker is shutting down")
default:
return fmt.Errorf("job queue is full")
}
}
func (w *Worker) GetJobStatus(jobID string) (*JobResult, bool) {
w.mu.RLock()
defer w.mu.RUnlock()
result, exists := w.results[jobID]
return result, exists
}
func (w *Worker) CancelJob(jobID string) error {
w.mu.Lock()
defer w.mu.Unlock()
if result, exists := w.results[jobID]; exists {
if result.Status == JobStatusRunning || result.Status == JobStatusPending {
result.Status = JobStatusCancelled
return nil
}
return fmt.Errorf("job cannot be cancelled")
}
return fmt.Errorf("job not found")
}
func (w *Worker) Shutdown() {
w.shuttingDown.Store(true)
w.cancel()
close(w.jobQueue)
w.wg.Wait()
}