diff --git a/cmd/silo/main.go b/cmd/silo/main.go index 48470fb1..07261062 100644 --- a/cmd/silo/main.go +++ b/cmd/silo/main.go @@ -2095,6 +2095,7 @@ func main() { } if ebookEnricher != nil { taskMgr.Register(tasks.NewSyncEbookMetadataTask(ebookEnricher)) + taskMgr.Register(tasks.NewBackfillEbookMetadataTask(ebookEnricher)) } if mangaEnricher != nil { taskMgr.Register(tasks.NewSyncMangaMetadataTask(mangaEnricher)) diff --git a/internal/ebooks/enrichment.go b/internal/ebooks/enrichment.go index 4cf6af36..900aefa4 100644 --- a/internal/ebooks/enrichment.go +++ b/internal/ebooks/enrichment.go @@ -99,7 +99,8 @@ type enrichmentItemRow struct { } type enrichmentQueue interface { - ClaimBatch(ctx context.Context, limit int, leaseDuration time.Duration) ([]EnrichmentJob, error) + ClaimBatch(ctx context.Context, scope EnrichmentScope, limit int, leaseDuration time.Duration) ([]EnrichmentJob, error) + ReadyCount(ctx context.Context, scope EnrichmentScope) (int, error) CheckClaim(ctx context.Context, job EnrichmentJob) error Complete(ctx context.Context, job EnrichmentJob, outcome EnrichmentOutcome, refreshAfter time.Duration) error Fail(ctx context.Context, job EnrichmentJob, errorClass EnrichmentErrorClass, message string, retryAfter time.Duration) error @@ -107,6 +108,15 @@ type enrichmentQueue interface { Discard(ctx context.Context, job EnrichmentJob) error } +type EnrichmentRunResult struct { + Claimed int `json:"claimed"` + Enriched int `json:"enriched"` + NoMatch int `json:"no_match"` + Failed int `json:"failed"` + Deferred int `json:"deferred"` + Remaining int `json:"remaining"` +} + // Enricher drives the ebook metadata enrichment sweep. type Enricher struct { pool *pgxpool.Pool @@ -173,9 +183,12 @@ func (e *Enricher) SetLiteraryWorkLinker(linker literaryWorkLinker) { e.workLinker = linker } -func (e *Enricher) Run(ctx context.Context) (int, error) { +func (e *Enricher) Run(ctx context.Context, scope EnrichmentScope) (EnrichmentRunResult, error) { if e == nil { - return 0, nil + return EnrichmentRunResult{}, nil + } + if err := scope.validate(); err != nil { + return EnrichmentRunResult{}, err } queue := e.queue @@ -183,14 +196,18 @@ func (e *Enricher) Run(ctx context.Context) (int, error) { queue = NewEnrichmentQueue(e.pool) } if queue == nil || (e.chainRepo == nil && e.enrichClaimedItemFn == nil) { - return 0, nil + return EnrichmentRunResult{}, nil } - jobs, err := queue.ClaimBatch(ctx, e.claimLimit(), defaultEnrichmentLease) + jobs, err := queue.ClaimBatch(ctx, scope, e.claimLimit(), defaultEnrichmentLease) if err != nil { - return 0, fmt.Errorf("ebook enrichment: claim batch: %w", err) + return EnrichmentRunResult{}, fmt.Errorf("ebook enrichment: claim batch: %w", err) } if len(jobs) == 0 { - return 0, nil + remaining, countErr := queue.ReadyCount(ctx, scope) + if countErr != nil { + return EnrichmentRunResult{}, fmt.Errorf("ebook enrichment: count remaining: %w", countErr) + } + return EnrichmentRunResult{Remaining: remaining}, nil } loadItems := e.loadClaimedItems @@ -200,7 +217,8 @@ func (e *Enricher) Run(ctx context.Context) (int, error) { items, err := loadItems(ctx, jobs) if err != nil { e.releaseJobs(queue, jobs) - return 0, fmt.Errorf("ebook enrichment: load claimed items: %w", err) + return EnrichmentRunResult{Claimed: len(jobs), Deferred: len(jobs)}, + fmt.Errorf("ebook enrichment: load claimed items: %w", err) } slog.InfoContext(ctx, "ebook enrichment: sweep started", "component", "ebooks", @@ -212,13 +230,22 @@ func (e *Enricher) Run(ctx context.Context) (int, error) { if e.enrichClaimedItemFn != nil { enrichItem = e.enrichClaimedItemFn } - enriched, runErr := e.runQueueBatch(ctx, queue, jobs, items, enrichItem) + result, runErr := e.runQueueBatch(ctx, queue, jobs, items, enrichItem) + remaining, countErr := queue.ReadyCount(ctx, scope) + result.Remaining = remaining + if countErr != nil { + runErr = errors.Join(runErr, fmt.Errorf("ebook enrichment: count remaining: %w", countErr)) + } slog.InfoContext(ctx, "ebook enrichment: sweep complete", "component", "ebooks", "attempted", len(items), - "enriched", enriched, + "enriched", result.Enriched, + "no_match", result.NoMatch, + "failed", result.Failed, + "deferred", result.Deferred, + "remaining", result.Remaining, ) - return enriched, runErr + return result, runErr } func (e *Enricher) claimLimit() int { @@ -239,7 +266,8 @@ func (e *Enricher) runQueueBatch( jobs []EnrichmentJob, items []enrichmentItemRow, enrichFn func(context.Context, enrichmentItemRow) (EnrichmentOutcome, error), -) (int, error) { +) (EnrichmentRunResult, error) { + result := EnrichmentRunResult{Claimed: len(jobs)} claimedJobs := make(map[string]EnrichmentJob, len(jobs)) for _, job := range jobs { claimedJobs[job.ContentID] = job @@ -253,6 +281,8 @@ func (e *Enricher) runQueueBatch( if _, ok := loaded[job.ContentID]; !ok { if err := e.discardJob(queue, job); err != nil && !errors.Is(err, ErrEnrichmentLeaseLost) { transitionErrs = append(transitionErrs, fmt.Errorf("%s: %w", job.ContentID, err)) + } else if err == nil { + result.Deferred++ } } } @@ -265,7 +295,7 @@ func (e *Enricher) runQueueBatch( workers = len(items) } if workers == 0 { - return 0, errors.Join(transitionErrs...) + return result, errors.Join(transitionErrs...) } itemTimeout := e.itemTimeout if itemTimeout <= 0 || itemTimeout >= defaultEnrichmentLease { @@ -276,6 +306,9 @@ func (e *Enricher) runQueueBatch( var ( wg sync.WaitGroup enriched int64 + noMatch int64 + failed int64 + deferred int64 transitionMu sync.Mutex ) recordTransitionError := func(contentID string, err error) { @@ -324,6 +357,9 @@ func (e *Enricher) runQueueBatch( 0, ) recordTransitionError(item.ContentID, transitionErr) + if transitionErr == nil { + atomic.AddInt64(&failed, 1) + } continue } if errors.Is(enrichErr, context.Canceled) { @@ -346,6 +382,9 @@ func (e *Enricher) runQueueBatch( recordTransitionError(item.ContentID, e.releaseJob(queue, job)) } recordTransitionError(item.ContentID, transitionErr) + if transitionErr == nil { + atomic.AddInt64(&failed, 1) + } continue } transitionErr := e.completeJob(queue, ctx, job, outcome) @@ -360,6 +399,10 @@ func (e *Enricher) runQueueBatch( } if outcome == EnrichmentOutcomeSuccess { atomic.AddInt64(&enriched, 1) + } else if outcome == EnrichmentOutcomeNoMatch { + atomic.AddInt64(&noMatch, 1) + } else { + atomic.AddInt64(&deferred, 1) } } }() @@ -373,7 +416,11 @@ func (e *Enricher) runQueueBatch( if ctx.Err() != nil { transitionErrs = append([]error{ctx.Err()}, transitionErrs...) } - return int(enriched), errors.Join(transitionErrs...) + result.Enriched += int(enriched) + result.NoMatch += int(noMatch) + result.Failed += int(failed) + result.Deferred += int(deferred) + return result, errors.Join(transitionErrs...) } func (e *Enricher) releaseJobs(queue enrichmentQueue, jobs []EnrichmentJob) { diff --git a/internal/ebooks/enrichment_queue.go b/internal/ebooks/enrichment_queue.go index 5cfefbac..5f92e5a6 100644 --- a/internal/ebooks/enrichment_queue.go +++ b/internal/ebooks/enrichment_queue.go @@ -37,6 +37,22 @@ const ( var ErrEnrichmentLeaseLost = errors.New("ebook enrichment lease lost") +type EnrichmentScope string + +const ( + EnrichmentScopeIncremental EnrichmentScope = "incremental" + EnrichmentScopeLegacy EnrichmentScope = "legacy" +) + +func (s EnrichmentScope) validate() error { + switch s { + case EnrichmentScopeIncremental, EnrichmentScopeLegacy: + return nil + default: + return fmt.Errorf("unsupported ebook enrichment scope %q", s) + } +} + type EnrichmentJob struct { ContentID string Token string @@ -144,6 +160,10 @@ var claimEnrichmentJobsQuery = ` FROM ebook_enrichment_state WHERE next_attempt_at <= now() AND (status = 'pending' OR (status = 'running' AND lease_until < now())) + AND ( + ($3 = 'incremental' AND priority >= 0) + OR ($3 = 'legacy' AND priority < 0) + ) ORDER BY (priority + FLOOR(EXTRACT(EPOCH FROM (now() - next_attempt_at)) / 3600)::integer) DESC, priority DESC, @@ -164,10 +184,18 @@ var claimEnrichmentJobsQuery = ` RETURNING state.content_id, state.claim_token, state.attempts, state.last_attempt_at, state.protected_fields ` -func (q *EnrichmentQueue) ClaimBatch(ctx context.Context, limit int, leaseDuration time.Duration) ([]EnrichmentJob, error) { +func (q *EnrichmentQueue) ClaimBatch( + ctx context.Context, + scope EnrichmentScope, + limit int, + leaseDuration time.Duration, +) ([]EnrichmentJob, error) { if q == nil || q.pool == nil { return nil, errors.New("ebook enrichment queue is not configured") } + if err := scope.validate(); err != nil { + return nil, err + } if limit <= 0 { return nil, nil } @@ -175,7 +203,7 @@ func (q *EnrichmentQueue) ClaimBatch(ctx context.Context, limit int, leaseDurati leaseDuration = defaultEnrichmentLease } - rows, err := q.pool.Query(ctx, claimEnrichmentJobsQuery, limit, postgresInterval(leaseDuration)) + rows, err := q.pool.Query(ctx, claimEnrichmentJobsQuery, limit, postgresInterval(leaseDuration), string(scope)) if err != nil { return nil, err } @@ -195,6 +223,31 @@ func (q *EnrichmentQueue) ClaimBatch(ctx context.Context, limit int, leaseDurati return jobs, nil } +var countReadyEnrichmentJobsQuery = ` + SELECT COUNT(*) + FROM ebook_enrichment_state + WHERE next_attempt_at <= now() + AND (status = 'pending' OR (status = 'running' AND lease_until < now())) + AND ( + ($1 = 'incremental' AND priority >= 0) + OR ($1 = 'legacy' AND priority < 0) + ) +` + +func (q *EnrichmentQueue) ReadyCount(ctx context.Context, scope EnrichmentScope) (int, error) { + if q == nil || q.pool == nil { + return 0, errors.New("ebook enrichment queue is not configured") + } + if err := scope.validate(); err != nil { + return 0, err + } + var count int + if err := q.pool.QueryRow(ctx, countReadyEnrichmentJobsQuery, string(scope)).Scan(&count); err != nil { + return 0, err + } + return count, nil +} + var checkEnrichmentClaimQuery = ` SELECT EXISTS ( SELECT 1 diff --git a/internal/ebooks/enrichment_queue_test.go b/internal/ebooks/enrichment_queue_test.go index f9a08e05..97427810 100644 --- a/internal/ebooks/enrichment_queue_test.go +++ b/internal/ebooks/enrichment_queue_test.go @@ -22,6 +22,8 @@ func TestEnrichmentQueueClaimQueryUsesAtomicLeasedClaims(t *testing.T) { "FOR UPDATE SKIP LOCKED", "status = 'pending' OR (status = 'running' AND lease_until < now())", "next_attempt_at <= now()", + "$3 = 'incremental' AND priority >= 0", + "$3 = 'legacy' AND priority < 0", "priority + FLOOR(EXTRACT(EPOCH FROM (now() - next_attempt_at)) / 3600)::integer", "SET status = 'running'", "lease_until = now() + $2::interval", @@ -36,6 +38,32 @@ func TestEnrichmentQueueClaimQueryUsesAtomicLeasedClaims(t *testing.T) { } } +func TestEnrichmentQueueReadyCountUsesTheSameScopeAsClaims(t *testing.T) { + query := strings.Join(strings.Fields(countReadyEnrichmentJobsQuery), " ") + for _, fragment := range []string{ + "SELECT COUNT(*)", + "next_attempt_at <= now()", + "status = 'pending' OR (status = 'running' AND lease_until < now())", + "$1 = 'incremental' AND priority >= 0", + "$1 = 'legacy' AND priority < 0", + } { + if !strings.Contains(query, fragment) { + t.Fatalf("ready-count query missing %q:\n%s", fragment, countReadyEnrichmentJobsQuery) + } + } +} + +func TestEnrichmentScopeValidation(t *testing.T) { + for _, scope := range []EnrichmentScope{EnrichmentScopeIncremental, EnrichmentScopeLegacy} { + if err := scope.validate(); err != nil { + t.Fatalf("validate(%q) error = %v", scope, err) + } + } + if err := EnrichmentScope("everything").validate(); err == nil { + t.Fatal("unknown enrichment scope was accepted") + } +} + func TestEnrichmentQueueCapturesEveryProviderWritablePendingField(t *testing.T) { query := strings.Join(strings.Fields(ebookProtectedFieldsSQL), " ") for _, field := range []string{ diff --git a/internal/ebooks/enrichment_test.go b/internal/ebooks/enrichment_test.go index 9f9afe4e..43b619e2 100644 --- a/internal/ebooks/enrichment_test.go +++ b/internal/ebooks/enrichment_test.go @@ -227,12 +227,13 @@ func TestEnricherRunTransitionsClaimedJobsByOutcome(t *testing.T) { }, } - enriched, err := e.Run(context.Background()) + result, err := e.Run(context.Background(), EnrichmentScopeIncremental) if err != nil { t.Fatalf("Run() error = %v", err) } - if enriched != 1 { - t.Fatalf("Run() enriched = %d, want 1", enriched) + want := EnrichmentRunResult{Claimed: 4, Enriched: 1, NoMatch: 1, Failed: 1, Deferred: 1} + if result != want { + t.Fatalf("Run() result = %+v, want %+v", result, want) } if queue.claimCalls != 1 { t.Fatalf("claim calls = %d, want 1", queue.claimCalls) @@ -281,7 +282,7 @@ func TestEnricherRunReleasesEveryLeaseOnCancellation(t *testing.T) { }, } - _, err := e.Run(ctx) + _, err := e.Run(ctx, EnrichmentScopeIncremental) if !errors.Is(err, context.Canceled) { t.Fatalf("Run() error = %v, want context.Canceled", err) } @@ -316,7 +317,7 @@ func TestEnricherRunReleasesLeaseWhenCompletionLosesContext(t *testing.T) { }, } - _, err := e.Run(ctx) + _, err := e.Run(ctx, EnrichmentScopeIncremental) if !errors.Is(err, context.Canceled) { t.Fatalf("Run() error = %v, want context.Canceled", err) } @@ -347,7 +348,7 @@ func TestEnricherRunDiscardsClaimedRowsThatAreNoLongerEligible(t *testing.T) { }, } - if _, err := e.Run(context.Background()); err != nil { + if _, err := e.Run(context.Background(), EnrichmentScopeIncremental); err != nil { t.Fatalf("Run() error = %v", err) } if got := strings.Join(queue.discarded, ","); got != "became-manga" { @@ -380,7 +381,7 @@ func TestEnricherRunBoundsEachItemBelowTheLease(t *testing.T) { } started := time.Now() - if _, err := e.Run(context.Background()); err != nil { + if _, err := e.Run(context.Background(), EnrichmentScopeIncremental); err != nil { t.Fatalf("Run() error = %v", err) } if elapsed := time.Since(started); elapsed > time.Second { @@ -408,7 +409,7 @@ func TestEnricherRunClaimsAtMostOneJobPerWorker(t *testing.T) { }, } - if _, err := e.Run(context.Background()); err != nil { + if _, err := e.Run(context.Background(), EnrichmentScopeIncremental); err != nil { t.Fatalf("Run() error = %v", err) } if queue.claimLimit != 4 { @@ -437,7 +438,7 @@ func TestEnricherRunInjectsExactClaimOwnershipCheck(t *testing.T) { }, } - if _, err := e.Run(context.Background()); err != nil { + if _, err := e.Run(context.Background(), EnrichmentScopeIncremental); err != nil { t.Fatalf("Run() error = %v", err) } if queue.claimCheckCalls != 1 { @@ -1084,6 +1085,8 @@ type fakeEnrichmentQueue struct { claimCalls int claimLimit int leaseDuration time.Duration + claimScope EnrichmentScope + remaining int completed map[string]EnrichmentOutcome failed map[string]EnrichmentErrorClass released []string @@ -1097,15 +1100,22 @@ type fakeEnrichmentQueue struct { completeErr error } -func (f *fakeEnrichmentQueue) ClaimBatch(_ context.Context, limit int, leaseDuration time.Duration) ([]EnrichmentJob, error) { +func (f *fakeEnrichmentQueue) ClaimBatch(_ context.Context, scope EnrichmentScope, limit int, leaseDuration time.Duration) ([]EnrichmentJob, error) { f.mu.Lock() defer f.mu.Unlock() f.claimCalls++ f.claimLimit = limit f.leaseDuration = leaseDuration + f.claimScope = scope return append([]EnrichmentJob(nil), f.jobs...), nil } +func (f *fakeEnrichmentQueue) ReadyCount(_ context.Context, _ EnrichmentScope) (int, error) { + f.mu.Lock() + defer f.mu.Unlock() + return f.remaining, nil +} + func (f *fakeEnrichmentQueue) Complete(_ context.Context, job EnrichmentJob, outcome EnrichmentOutcome, _ time.Duration) error { f.mu.Lock() defer f.mu.Unlock() diff --git a/internal/taskmanager/tasks/sync_ebook_metadata.go b/internal/taskmanager/tasks/sync_ebook_metadata.go index 8656ab7d..082c2b47 100644 --- a/internal/taskmanager/tasks/sync_ebook_metadata.go +++ b/internal/taskmanager/tasks/sync_ebook_metadata.go @@ -4,53 +4,157 @@ import ( "context" "encoding/json" "fmt" + "time" + "github.com/Silo-Server/silo-server/internal/ebooks" "github.com/Silo-Server/silo-server/internal/taskmanager" ) +const ebookMetadataExecutionBudget = 4 * time.Minute + type ebookMetadataEnricher interface { - Run(ctx context.Context) (int, error) + Run(ctx context.Context, scope ebooks.EnrichmentScope) (ebooks.EnrichmentRunResult, error) } -// SyncEbookMetadataTask runs the periodic ebook enrichment sweep. -// It calls ebooks.Enricher.Run() which selects unenriched ebook media_items, -// resolves the per-folder metadata-provider chain at content_level='ebook', -// and writes results back to the database. +type ebookMetadataTask struct { + enricher ebookMetadataEnricher + scope ebooks.EnrichmentScope + key string + name string + description string + triggers []taskmanager.TriggerConfig + budget time.Duration + now func() time.Time + errorPrefix string +} + +// SyncEbookMetadataTask drains new and recurring ebook metadata work. type SyncEbookMetadataTask struct { - enricher ebookMetadataEnricher + *ebookMetadataTask +} + +// BackfillEbookMetadataTask drains the legacy ebook backlog only when manually run. +type BackfillEbookMetadataTask struct { + *ebookMetadataTask } -// NewSyncEbookMetadataTask constructs the task. func NewSyncEbookMetadataTask(enricher ebookMetadataEnricher) *SyncEbookMetadataTask { - return &SyncEbookMetadataTask{enricher: enricher} + return &SyncEbookMetadataTask{ebookMetadataTask: &ebookMetadataTask{ + enricher: enricher, + scope: ebooks.EnrichmentScopeIncremental, + key: "sync_ebook_metadata", + name: "Sync Ebook Metadata", + description: "Fetches metadata for new ebooks and ebooks due for a recurring refresh", + triggers: []taskmanager.TriggerConfig{{Type: taskmanager.TriggerTypeInterval, IntervalMs: 5 * 60 * 1000}}, + budget: ebookMetadataExecutionBudget, + now: time.Now, + errorPrefix: "ebook metadata sync", + }} } -func (t *SyncEbookMetadataTask) Key() string { return "sync_ebook_metadata" } -func (t *SyncEbookMetadataTask) Name() string { return "Sync Ebook Metadata" } -func (t *SyncEbookMetadataTask) Description() string { - return "Fetches metadata (cover art, overview, authors) for ebooks that have not yet been enriched" +func NewBackfillEbookMetadataTask(enricher ebookMetadataEnricher) *BackfillEbookMetadataTask { + return &BackfillEbookMetadataTask{ebookMetadataTask: &ebookMetadataTask{ + enricher: enricher, + scope: ebooks.EnrichmentScopeLegacy, + key: "backfill_ebook_metadata", + name: "Backfill Ebook Metadata", + description: "Manually enriches the legacy ebook backlog without competing with scheduled metadata sync", + budget: ebookMetadataExecutionBudget, + now: time.Now, + errorPrefix: "ebook metadata backfill", + }} } -func (t *SyncEbookMetadataTask) Category() taskmanager.TaskCategory { + +func (t *ebookMetadataTask) Key() string { return t.key } +func (t *ebookMetadataTask) Name() string { return t.name } +func (t *ebookMetadataTask) Description() string { return t.description } +func (t *ebookMetadataTask) Category() taskmanager.TaskCategory { return taskmanager.TaskCategoryMetadata } -func (t *SyncEbookMetadataTask) IsHidden() bool { return false } +func (t *ebookMetadataTask) IsHidden() bool { return false } -func (t *SyncEbookMetadataTask) DefaultTriggers() []taskmanager.TriggerConfig { - return []taskmanager.TriggerConfig{ - {Type: taskmanager.TriggerTypeInterval, IntervalMs: 5 * 60 * 1000}, +func (t *ebookMetadataTask) DefaultTriggers() []taskmanager.TriggerConfig { + return append([]taskmanager.TriggerConfig(nil), t.triggers...) +} + +func (t *ebookMetadataTask) Execute(ctx context.Context, progress taskmanager.ProgressReporter) error { + progress.Report(0, fmt.Sprintf("%s started", t.name)) + started := t.now() + total := ebooks.EnrichmentRunResult{} + + for { + if err := ctx.Err(); err != nil { + return err + } + + batch, err := t.enricher.Run(ctx, t.scope) + if err != nil { + return fmt.Errorf("%s: %w", t.errorPrefix, err) + } + addEbookEnrichmentResult(&total, batch) + reportEbookEnrichmentProgress(progress, total, false) + + if err := ctx.Err(); err != nil { + return err + } + if batch.Remaining == 0 { + reportEbookEnrichmentProgress(progress, total, true) + return nil + } + if batch.Claimed == 0 { + return nil + } + if t.now().Sub(started) >= t.budget { + total.Deferred += total.Remaining + reportEbookEnrichmentProgress(progress, total, false) + return nil + } } } -func (t *SyncEbookMetadataTask) Execute(ctx context.Context, progress taskmanager.ProgressReporter) error { - progress.Report(0, "Scanning for unenriched ebooks") - - enriched, err := t.enricher.Run(ctx) - if err != nil { - return fmt.Errorf("ebook metadata sync: %w", err) - } - - result, _ := json.Marshal(map[string]int{"items_enriched": enriched}) - progress.SetResultData(result) - progress.Report(100, fmt.Sprintf("Ebook metadata sync complete (%d items enriched)", enriched)) - return nil +func addEbookEnrichmentResult(total *ebooks.EnrichmentRunResult, batch ebooks.EnrichmentRunResult) { + total.Claimed += batch.Claimed + total.Enriched += batch.Enriched + total.NoMatch += batch.NoMatch + total.Failed += batch.Failed + total.Deferred += batch.Deferred + total.Remaining = batch.Remaining +} + +func reportEbookEnrichmentProgress( + progress taskmanager.ProgressReporter, + result ebooks.EnrichmentRunResult, + complete bool, +) { + data, _ := json.Marshal(result) + progress.SetResultData(data) + + percent := ebookEnrichmentPercent(result) + if complete && result.Remaining == 0 { + percent = 100 + } + progress.Report(percent, fmt.Sprintf( + "Claimed %d, enriched %d, no match %d, failed %d, deferred %d, remaining %d", + result.Claimed, + result.Enriched, + result.NoMatch, + result.Failed, + result.Deferred, + result.Remaining, + )) +} + +func ebookEnrichmentPercent(result ebooks.EnrichmentRunResult) float64 { + if result.Remaining == 0 { + return 100 + } + total := result.Claimed + result.Remaining + if total <= 0 { + return 0 + } + percent := float64(result.Claimed) * 100 / float64(total) + if percent >= 100 { + return 99 + } + return percent } diff --git a/internal/taskmanager/tasks/sync_ebook_metadata_test.go b/internal/taskmanager/tasks/sync_ebook_metadata_test.go index 0adb8c9a..75ea8806 100644 --- a/internal/taskmanager/tasks/sync_ebook_metadata_test.go +++ b/internal/taskmanager/tasks/sync_ebook_metadata_test.go @@ -6,86 +6,174 @@ import ( "errors" "strings" "testing" + "time" + "github.com/Silo-Server/silo-server/internal/ebooks" "github.com/Silo-Server/silo-server/internal/taskmanager" ) type fakeEbookMetadataEnricher struct { - enriched int - err error - called bool + results []ebooks.EnrichmentRunResult + err error + scopes []ebooks.EnrichmentScope + onRun func(int) } -func (f *fakeEbookMetadataEnricher) Run(context.Context) (int, error) { - f.called = true - return f.enriched, f.err +func (f *fakeEbookMetadataEnricher) Run(_ context.Context, scope ebooks.EnrichmentScope) (ebooks.EnrichmentRunResult, error) { + f.scopes = append(f.scopes, scope) + call := len(f.scopes) + if f.onRun != nil { + f.onRun(call) + } + if f.err != nil { + return ebooks.EnrichmentRunResult{}, f.err + } + if call > len(f.results) { + return ebooks.EnrichmentRunResult{}, nil + } + return f.results[call-1], nil } type ebookMetadataProgressReporter struct { + percents []float64 messages []string - result json.RawMessage + results []json.RawMessage } -func (p *ebookMetadataProgressReporter) Report(_ float64, message string) { +func (p *ebookMetadataProgressReporter) Report(percent float64, message string) { + p.percents = append(p.percents, percent) p.messages = append(p.messages, message) } func (p *ebookMetadataProgressReporter) SetResultData(data json.RawMessage) { - p.result = data + p.results = append(p.results, append(json.RawMessage(nil), data...)) } -func TestSyncEbookMetadataTaskProperties(t *testing.T) { - task := NewSyncEbookMetadataTask(&fakeEbookMetadataEnricher{}) +func TestEbookMetadataTaskPropertiesAndScopes(t *testing.T) { + enricher := &fakeEbookMetadataEnricher{} + syncTask := NewSyncEbookMetadataTask(enricher) + backfillTask := NewBackfillEbookMetadataTask(enricher) - if task.Key() != "sync_ebook_metadata" { - t.Fatalf("Key() = %q, want sync_ebook_metadata", task.Key()) + if syncTask.Key() != "sync_ebook_metadata" || syncTask.Name() != "Sync Ebook Metadata" { + t.Fatalf("unexpected sync identity: %q %q", syncTask.Key(), syncTask.Name()) } - if task.Name() != "Sync Ebook Metadata" { - t.Fatalf("Name() = %q, want Sync Ebook Metadata", task.Name()) + if backfillTask.Key() != "backfill_ebook_metadata" { + t.Fatalf("backfill Key() = %q", backfillTask.Key()) } - if task.Category() != taskmanager.TaskCategoryMetadata { - t.Fatalf("Category() = %q, want %q", task.Category(), taskmanager.TaskCategoryMetadata) + if !strings.Contains(strings.ToLower(backfillTask.Description()), "legacy") { + t.Fatalf("backfill description does not explain legacy work: %q", backfillTask.Description()) } - if task.IsHidden() { - t.Fatal("IsHidden() = true, want false") + for _, task := range []taskmanager.Task{syncTask, backfillTask} { + if task.Category() != taskmanager.TaskCategoryMetadata || task.IsHidden() { + t.Fatalf("unexpected task properties for %q", task.Key()) + } } - triggers := task.DefaultTriggers() + triggers := syncTask.DefaultTriggers() if len(triggers) != 1 || triggers[0].Type != taskmanager.TriggerTypeInterval || triggers[0].IntervalMs != 5*60*1000 { - t.Fatalf("DefaultTriggers() = %#v, want one 5 minute interval", triggers) + t.Fatalf("sync DefaultTriggers() = %#v", triggers) } - if strings.Contains(strings.ToLower(task.Description()), "narrator") { - t.Fatalf("Description() mentions narrator: %q", task.Description()) + if triggers := backfillTask.DefaultTriggers(); len(triggers) != 0 { + t.Fatalf("backfill DefaultTriggers() = %#v, want none", triggers) } } -func TestSyncEbookMetadataTaskExecuteReportsResult(t *testing.T) { - enricher := &fakeEbookMetadataEnricher{enriched: 3} +func TestEbookMetadataTaskDrainsBatchesAndReportsHonestProgress(t *testing.T) { + enricher := &fakeEbookMetadataEnricher{results: []ebooks.EnrichmentRunResult{ + {Claimed: 4, Enriched: 2, NoMatch: 1, Failed: 1, Remaining: 3}, + {Claimed: 3, Enriched: 1, Deferred: 2, Remaining: 0}, + }} task := NewSyncEbookMetadataTask(enricher) progress := &ebookMetadataProgressReporter{} if err := task.Execute(context.Background(), progress); err != nil { t.Fatalf("Execute() error = %v", err) } - if !enricher.called { - t.Fatal("enricher was not called") + if len(enricher.scopes) != 2 { + t.Fatalf("Run calls = %d, want 2", len(enricher.scopes)) } - var result map[string]int - if err := json.Unmarshal(progress.result, &result); err != nil { - t.Fatalf("result data JSON error: %v", err) + for _, scope := range enricher.scopes { + if scope != ebooks.EnrichmentScopeIncremental { + t.Fatalf("sync scope = %q, want incremental", scope) + } } - if result["items_enriched"] != 3 { - t.Fatalf("items_enriched = %d, want 3", result["items_enriched"]) + var result ebooks.EnrichmentRunResult + if err := json.Unmarshal(progress.results[len(progress.results)-1], &result); err != nil { + t.Fatalf("result JSON error: %v", err) } - if len(progress.messages) == 0 || !strings.Contains(progress.messages[len(progress.messages)-1], "3 items enriched") { - t.Fatalf("progress messages = %#v, want completion count", progress.messages) + want := ebooks.EnrichmentRunResult{Claimed: 7, Enriched: 3, NoMatch: 1, Failed: 1, Deferred: 2} + if result != want { + t.Fatalf("result = %+v, want %+v", result, want) + } + if progress.percents[1] >= 100 { + t.Fatalf("first batch progress = %.1f, must be below 100 with remaining work", progress.percents[1]) + } + if got := progress.percents[len(progress.percents)-1]; got != 100 { + t.Fatalf("final progress = %.1f, want 100", got) } } -func TestSyncEbookMetadataTaskExecuteWrapsError(t *testing.T) { - enricher := &fakeEbookMetadataEnricher{err: errors.New("boom")} - task := NewSyncEbookMetadataTask(enricher) +func TestEbookMetadataBackfillUsesLegacyScope(t *testing.T) { + enricher := &fakeEbookMetadataEnricher{results: []ebooks.EnrichmentRunResult{{Remaining: 0}}} + task := NewBackfillEbookMetadataTask(enricher) + if err := task.Execute(context.Background(), &ebookMetadataProgressReporter{}); err != nil { + t.Fatalf("Execute() error = %v", err) + } + if len(enricher.scopes) != 1 || enricher.scopes[0] != ebooks.EnrichmentScopeLegacy { + t.Fatalf("backfill scopes = %#v, want legacy", enricher.scopes) + } +} - err := task.Execute(context.Background(), &ebookMetadataProgressReporter{}) +func TestEbookMetadataTaskStopsBetweenBatchesOnCancellation(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + enricher := &fakeEbookMetadataEnricher{ + results: []ebooks.EnrichmentRunResult{{Claimed: 1, Enriched: 1, Remaining: 2}}, + onRun: func(int) { + cancel() + }, + } + err := NewSyncEbookMetadataTask(enricher).Execute(ctx, &ebookMetadataProgressReporter{}) + if !errors.Is(err, context.Canceled) { + t.Fatalf("Execute() error = %v, want context.Canceled", err) + } + if len(enricher.scopes) != 1 { + t.Fatalf("Run calls = %d, want 1", len(enricher.scopes)) + } +} + +func TestEbookMetadataTaskTreatsBudgetAsCleanDeferredCompletion(t *testing.T) { + now := time.Date(2026, 7, 19, 12, 0, 0, 0, time.UTC) + enricher := &fakeEbookMetadataEnricher{ + results: []ebooks.EnrichmentRunResult{{Claimed: 2, Enriched: 2, Remaining: 5}}, + onRun: func(int) { + now = now.Add(5 * time.Minute) + }, + } + task := NewSyncEbookMetadataTask(enricher) + task.now = func() time.Time { return now } + task.budget = 4 * time.Minute + progress := &ebookMetadataProgressReporter{} + + if err := task.Execute(context.Background(), progress); err != nil { + t.Fatalf("Execute() error = %v", err) + } + if len(enricher.scopes) != 1 { + t.Fatalf("Run calls = %d, want 1", len(enricher.scopes)) + } + var result ebooks.EnrichmentRunResult + if err := json.Unmarshal(progress.results[len(progress.results)-1], &result); err != nil { + t.Fatalf("result JSON error: %v", err) + } + if result.Remaining != 5 || result.Deferred != 5 { + t.Fatalf("budget result = %+v, want remaining=5 deferred=5", result) + } + if got := progress.percents[len(progress.percents)-1]; got >= 100 { + t.Fatalf("budget-limited progress = %.1f, must remain below 100", got) + } +} + +func TestEbookMetadataTaskWrapsRunError(t *testing.T) { + enricher := &fakeEbookMetadataEnricher{err: errors.New("boom")} + err := NewSyncEbookMetadataTask(enricher).Execute(context.Background(), &ebookMetadataProgressReporter{}) if err == nil || !strings.Contains(err.Error(), "ebook metadata sync") { t.Fatalf("Execute() error = %v, want wrapped ebook metadata sync error", err) } diff --git a/migrations/sql/20260719090000_ebook_enrichment_jobs.sql b/migrations/sql/20260719090000_ebook_enrichment_jobs.sql index 092593df..535bf734 100644 --- a/migrations/sql/20260719090000_ebook_enrichment_jobs.sql +++ b/migrations/sql/20260719090000_ebook_enrichment_jobs.sql @@ -155,7 +155,7 @@ $$; -- +goose StatementEnd CREATE INDEX CONCURRENTLY IF NOT EXISTS ebook_enrichment_state_claim_idx - ON public.ebook_enrichment_state (next_attempt_at, priority DESC, updated_at) + ON public.ebook_enrichment_state (priority, next_attempt_at, updated_at) WHERE status IN ('pending', 'running'); -- +goose Down