diff --git a/internal/ebooks/enrichment.go b/internal/ebooks/enrichment.go index e6d9062e..02130497 100644 --- a/internal/ebooks/enrichment.go +++ b/internal/ebooks/enrichment.go @@ -19,6 +19,9 @@ import ( "time" "github.com/jackc/pgx/v5/pgxpool" + "google.golang.org/genproto/googleapis/rpc/errdetails" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" "github.com/Silo-Server/silo-server/internal/catalog" "github.com/Silo-Server/silo-server/internal/metadata" @@ -30,19 +33,11 @@ const ( defaultEnrichBatchSize = 50 defaultEnrichWorkers = 4 - - // enrichFailureCap is the ebook_enrichment_state.failures count at which - // an ebook stops being claimed for enrichment. Combined with the - // failure-count-first claim ordering this prevents a head-of-line block - // of permanently failing items from starving newer items and hammering - // providers. - enrichFailureCap = 5 ) -// errEnrichmentSkipped marks an item that could not be attempted at all (no -// library folder linked yet, no providers configured). Skipped items are -// neither stamped as refreshed nor counted against the failure cap, so they -// are retried on every sweep until the missing prerequisite appears. +// errEnrichmentSkipped preserves the direct helper contract for an item that +// cannot be attempted. Queue-backed runs record a short skipped horizon instead +// of counting the missing prerequisite as a provider failure. var errEnrichmentSkipped = errors.New("ebook enrichment skipped") func ebookContentType() string { @@ -70,6 +65,26 @@ type enrichmentItemRow struct { Language string Author string ProviderIDs map[string]string + Status string + Overview string + ReleaseDate string + Genres []string + Studios []string + PosterPath string +} + +type enrichmentQueue interface { + MaterializeCandidates(ctx context.Context) error + ClaimBatch(ctx context.Context, limit int, leaseDuration time.Duration) ([]EnrichmentJob, error) + Complete(ctx context.Context, contentID string, outcome EnrichmentOutcome, refreshAfter time.Duration) error + Fail(ctx context.Context, contentID string, errorClass EnrichmentErrorClass, message string, retryAfter time.Duration) error + Release(ctx context.Context, contentID string) error +} + +type claimAwareEnrichmentQueue interface { + CompleteClaim(ctx context.Context, job EnrichmentJob, outcome EnrichmentOutcome, refreshAfter time.Duration) error + FailClaim(ctx context.Context, job EnrichmentJob, errorClass EnrichmentErrorClass, message string, retryAfter time.Duration) error + ReleaseClaim(ctx context.Context, job EnrichmentJob) error } // Enricher drives the ebook metadata enrichment sweep. @@ -85,6 +100,10 @@ type Enricher struct { workLinker literaryWorkLinker batchSize int workers int + queue enrichmentQueue + + loadClaimedItemsFn func(context.Context, []EnrichmentJob) ([]enrichmentItemRow, error) + enrichClaimedItemFn func(context.Context, enrichmentItemRow) (EnrichmentOutcome, error) } type literaryWorkLinker interface { @@ -108,6 +127,7 @@ func NewEnricher( providerIDs: providerIDs, batchSize: defaultEnrichBatchSize, workers: ebookEnrichWorkers(), + queue: NewEnrichmentQueue(pool), } } @@ -133,30 +153,210 @@ func (e *Enricher) SetLiteraryWorkLinker(linker literaryWorkLinker) { } func (e *Enricher) Run(ctx context.Context) (int, error) { - if e == nil || e.pool == nil || e.chainRepo == nil { + if e == nil { return 0, nil } - items, err := e.claimBatch(ctx) + queue := e.queue + if queue == nil && e.pool != nil { + queue = NewEnrichmentQueue(e.pool) + } + if queue == nil || (e.chainRepo == nil && e.enrichClaimedItemFn == nil) { + return 0, nil + } + if err := queue.MaterializeCandidates(ctx); err != nil { + return 0, fmt.Errorf("ebook enrichment: materialize candidates: %w", err) + } + jobs, err := queue.ClaimBatch(ctx, e.batchSize, defaultEnrichmentLease) if err != nil { return 0, fmt.Errorf("ebook enrichment: claim batch: %w", err) } - if len(items) == 0 { + if len(jobs) == 0 { return 0, nil } + loadItems := e.loadClaimedItems + if e.loadClaimedItemsFn != nil { + loadItems = e.loadClaimedItemsFn + } + items, err := loadItems(ctx, jobs) + if err != nil { + e.releaseJobs(queue, jobs) + return 0, fmt.Errorf("ebook enrichment: load claimed items: %w", err) + } + slog.InfoContext(ctx, "ebook enrichment: sweep started", "component", "ebooks", "count", len(items), "workers", e.workers, ) - enriched := e.runBatch(ctx, items, e.enrichItem, e.recordEnrichFailure) + enrichItem := e.enrichClaimedItem + if e.enrichClaimedItemFn != nil { + enrichItem = e.enrichClaimedItemFn + } + enriched, runErr := e.runQueueBatch(ctx, queue, jobs, items, enrichItem) slog.InfoContext(ctx, "ebook enrichment: sweep complete", "component", "ebooks", "attempted", len(items), "enriched", enriched, ) - return enriched, nil + return enriched, runErr +} + +func (e *Enricher) runQueueBatch( + ctx context.Context, + queue enrichmentQueue, + jobs []EnrichmentJob, + items []enrichmentItemRow, + enrichFn func(context.Context, enrichmentItemRow) (EnrichmentOutcome, error), +) (int, error) { + claimedJobs := make(map[string]EnrichmentJob, len(jobs)) + for _, job := range jobs { + claimedJobs[job.ContentID] = job + } + loaded := make(map[string]struct{}, len(items)) + for i := range items { + loaded[items[i].ContentID] = struct{}{} + } + for _, job := range jobs { + if _, ok := loaded[job.ContentID]; !ok { + _ = e.releaseJob(queue, job) + } + } + + workers := e.workers + if workers <= 0 { + workers = 1 + } + if workers > len(items) { + workers = len(items) + } + if workers == 0 { + return 0, nil + } + + ch := make(chan enrichmentItemRow, workers) + var ( + wg sync.WaitGroup + enriched int64 + transitionMu sync.Mutex + transitionErrs []error + ) + recordTransitionError := func(contentID string, err error) { + if err == nil || errors.Is(err, ErrEnrichmentLeaseLost) { + return + } + transitionMu.Lock() + defer transitionMu.Unlock() + transitionErrs = append(transitionErrs, fmt.Errorf("%s: %w", contentID, err)) + } + + for i := 0; i < workers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for item := range ch { + job := claimedJobs[item.ContentID] + if ctx.Err() != nil { + recordTransitionError(item.ContentID, e.releaseJob(queue, job)) + continue + } + + outcome, enrichErr := enrichFn(ctx, item) + if ctx.Err() != nil || errors.Is(enrichErr, context.Canceled) || errors.Is(enrichErr, context.DeadlineExceeded) { + recordTransitionError(item.ContentID, e.releaseJob(queue, job)) + continue + } + if enrichErr != nil { + errorClass, retryAfter := classifyEnrichmentError(enrichErr) + transitionErr := e.failJob( + queue, + ctx, + job, + errorClass, + enrichErr.Error(), + retryAfter, + ) + if transitionErr != nil && (ctx.Err() != nil || + errors.Is(transitionErr, context.Canceled) || + errors.Is(transitionErr, context.DeadlineExceeded)) { + recordTransitionError(item.ContentID, e.releaseJob(queue, job)) + } + recordTransitionError(item.ContentID, transitionErr) + continue + } + transitionErr := e.completeJob(queue, ctx, job, outcome) + if transitionErr != nil && (ctx.Err() != nil || + errors.Is(transitionErr, context.Canceled) || + errors.Is(transitionErr, context.DeadlineExceeded)) { + recordTransitionError(item.ContentID, e.releaseJob(queue, job)) + } + if transitionErr != nil { + recordTransitionError(item.ContentID, transitionErr) + continue + } + if outcome == EnrichmentOutcomeSuccess { + atomic.AddInt64(&enriched, 1) + } + } + }() + } + for _, item := range items { + ch <- item + } + close(ch) + wg.Wait() + + if ctx.Err() != nil { + transitionErrs = append([]error{ctx.Err()}, transitionErrs...) + } + return int(enriched), errors.Join(transitionErrs...) +} + +func (e *Enricher) releaseJobs(queue enrichmentQueue, jobs []EnrichmentJob) { + for _, job := range jobs { + if err := e.releaseJob(queue, job); err != nil && !errors.Is(err, ErrEnrichmentLeaseLost) { + slog.Warn("ebook enrichment: failed to release lease", "component", "ebooks", + "content_id", job.ContentID, + "error", err, + ) + } + } +} + +func (e *Enricher) completeJob( + queue enrichmentQueue, + ctx context.Context, + job EnrichmentJob, + outcome EnrichmentOutcome, +) error { + if claimQueue, ok := queue.(claimAwareEnrichmentQueue); ok { + return claimQueue.CompleteClaim(ctx, job, outcome, enrichmentRefreshHorizon(outcome)) + } + return queue.Complete(ctx, job.ContentID, outcome, enrichmentRefreshHorizon(outcome)) +} + +func (e *Enricher) failJob( + queue enrichmentQueue, + ctx context.Context, + job EnrichmentJob, + errorClass EnrichmentErrorClass, + message string, + retryAfter time.Duration, +) error { + if claimQueue, ok := queue.(claimAwareEnrichmentQueue); ok { + return claimQueue.FailClaim(ctx, job, errorClass, message, retryAfter) + } + return queue.Fail(ctx, job.ContentID, errorClass, message, retryAfter) +} + +func (e *Enricher) releaseJob(queue enrichmentQueue, job EnrichmentJob) error { + releaseCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if claimQueue, ok := queue.(claimAwareEnrichmentQueue); ok { + return claimQueue.ReleaseClaim(releaseCtx, job) + } + return queue.Release(releaseCtx, job.ContentID) } func (e *Enricher) runBatch( @@ -222,15 +422,12 @@ func (e *Enricher) runBatch( return int(enriched) } -// claimBatchQuery selects unenriched ebooks. Items with fewer prior failures -// are claimed first and items at/above enrichFailureCap are skipped entirely, -// so a block of permanently failing items cannot occupy every sweep. -var claimBatchQuery = ` +var loadEnrichmentItemsQuery = ` SELECT mi.content_id, mi.title, - mi.year, - COALESCE(mil.media_folder_id, 0) AS folder_id, + COALESCE(mi.year, 0), + COALESCE(membership.media_folder_id, 0) AS folder_id, COALESCE(mf.metadata_language, 'en') AS language, COALESCE( (SELECT p.name @@ -241,33 +438,40 @@ var claimBatchQuery = ` ORDER BY ip.sort_order, ip.id LIMIT 1), '' - ) AS author - FROM media_items mi - LEFT JOIN media_item_libraries mil ON mil.content_id = mi.content_id - LEFT JOIN media_folders mf ON mf.id = mil.media_folder_id - LEFT JOIN ebook_enrichment_state ees ON ees.content_id = mi.content_id + ) AS author, + COALESCE(mi.status, ''), + COALESCE(mi.overview, ''), + COALESCE(mi.release_date, ''), + COALESCE(mi.genres, '{}'), + COALESCE(mi.studios, '{}'), + COALESCE(mi.poster_path, '') + FROM unnest($1::text[]) WITH ORDINALITY AS claimed(content_id, position) + JOIN media_items mi ON mi.content_id = claimed.content_id + LEFT JOIN LATERAL ( + SELECT mil.media_folder_id + FROM media_item_libraries mil + WHERE mil.content_id = mi.content_id + ORDER BY mil.first_seen_at, mil.media_folder_id + LIMIT 1 + ) membership ON true + LEFT JOIN media_folders mf ON mf.id = membership.media_folder_id WHERE mi.type = 'ebook' - -- Manga chapters are type='ebook' but are parts of a series, not - -- standalone books. They are enriched via their type='manga' series (a - -- separate path), never individually against book sources — excluding - -- them here stops a pointless search storm over Gutenberg/Anna's/etc. AND ` + catalog.MangaChapterExclusionWhere("mi") + ` - AND (mi.poster_path IS NULL OR mi.poster_path = '') - AND mi.last_refreshed IS NULL - AND COALESCE(ees.failures, 0) < $2 - ORDER BY COALESCE(ees.failures, 0) ASC, mi.created_at ASC - LIMIT $1 + ORDER BY claimed.position ` -func (e *Enricher) claimBatch(ctx context.Context) ([]enrichmentItemRow, error) { - rows, err := e.pool.Query(ctx, claimBatchQuery, e.batchSize, enrichFailureCap) +func (e *Enricher) loadClaimedItems(ctx context.Context, jobs []EnrichmentJob) ([]enrichmentItemRow, error) { + contentIDs := make([]string, 0, len(jobs)) + for _, job := range jobs { + contentIDs = append(contentIDs, job.ContentID) + } + rows, err := e.pool.Query(ctx, loadEnrichmentItemsQuery, contentIDs) if err != nil { - return nil, fmt.Errorf("querying unenriched ebooks: %w", err) + return nil, fmt.Errorf("querying claimed ebooks: %w", err) } defer rows.Close() var items []enrichmentItemRow - seen := make(map[string]struct{}) for rows.Next() { var item enrichmentItemRow if err := rows.Scan( @@ -277,13 +481,15 @@ func (e *Enricher) claimBatch(ctx context.Context) ([]enrichmentItemRow, error) &item.FolderID, &item.Language, &item.Author, + &item.Status, + &item.Overview, + &item.ReleaseDate, + &item.Genres, + &item.Studios, + &item.PosterPath, ); err != nil { return nil, fmt.Errorf("scanning ebook enrichment row: %w", err) } - if _, dup := seen[item.ContentID]; dup { - continue - } - seen[item.ContentID] = struct{}{} items = append(items, item) } if err := rows.Err(); err != nil { @@ -303,19 +509,27 @@ func (e *Enricher) claimBatch(ctx context.Context) ([]enrichmentItemRow, error) } func (e *Enricher) enrichItem(ctx context.Context, item enrichmentItemRow) error { + outcome, err := e.enrichClaimedItem(ctx, item) + if err == nil && outcome == EnrichmentOutcomeSkipped { + return fmt.Errorf("%w: item %s is not ready", errEnrichmentSkipped, item.ContentID) + } + return err +} + +func (e *Enricher) enrichClaimedItem(ctx context.Context, item enrichmentItemRow) (EnrichmentOutcome, error) { if item.FolderID == 0 { // The scanner inserts the library membership after the item upsert, so // a freshly indexed ebook can be claimed inside that window. Skip it: // stamping here would terminally mark the item refreshed before any // provider ever saw it. - return fmt.Errorf("%w: item %s has no library folder yet", errEnrichmentSkipped, item.ContentID) + return EnrichmentOutcomeSkipped, nil } providers, err := metadata.ResolveChain(ctx, item.FolderID, ebookContentType(), e.chainRepo, e.resolver) if err != nil { - return fmt.Errorf("resolving ebook chain for folder %d: %w", item.FolderID, err) + return "", fmt.Errorf("resolving ebook chain for folder %d: %w", item.FolderID, err) } - return e.enrichWithProviders(ctx, item, providers) + return e.enrichWithProvidersOutcome(ctx, item, providers) } // enrichWithProviders runs the provider chain for one claimed item. Outcomes: @@ -323,13 +537,25 @@ func (e *Enricher) enrichItem(ctx context.Context, item enrichmentItemRow) error // - providers answered but nothing matched: stamp last_refreshed so the // item is not re-claimed every sweep (nil error); // - one or more providers errored and no metadata was obtained: return an -// error so the failure cap/backoff engages, without stamping; +// error so durable queue backoff engages, without stamping; // - no providers configured: skip (no stamp, no failure) so the item is // retried once a chain exists. func (e *Enricher) enrichWithProviders(ctx context.Context, item enrichmentItemRow, providers []metadata.Provider) error { - if len(providers) == 0 { + outcome, err := e.enrichWithProvidersOutcome(ctx, item, providers) + if err == nil && outcome == EnrichmentOutcomeSkipped { return fmt.Errorf("%w: no metadata providers configured for folder %d", errEnrichmentSkipped, item.FolderID) } + return err +} + +func (e *Enricher) enrichWithProvidersOutcome( + ctx context.Context, + item enrichmentItemRow, + providers []metadata.Provider, +) (EnrichmentOutcome, error) { + if len(providers) == 0 { + return EnrichmentOutcomeSkipped, nil + } var owner providerIDOwnerLookup if e.providerIDs != nil { @@ -340,23 +566,25 @@ func (e *Enricher) enrichWithProviders(ctx context.Context, item enrichmentItemR if !accumulator.HasMetadata && accumulator.PosterPath == "" && accumulator.Overview == "" { if err := ctx.Err(); err != nil { // A cancelled sweep says nothing about the item or the providers. - return err + return "", err } if len(providerErrs) > 0 { - // Transient provider trouble must not stamp the item terminally; - // surfacing an error engages the failure cap and backoff instead. - return fmt.Errorf("no metadata obtained, %d provider error(s): %w", + return "", fmt.Errorf("no metadata obtained, %d provider error(s): %w", len(providerErrs), errors.Join(providerErrs...)) } slog.InfoContext(ctx, "ebook enrichment: no metadata found", "component", "ebooks", "content_id", item.ContentID, "title", item.Title, ) - return e.stampLastRefreshed(ctx, item.ContentID) + if err := e.stampLastRefreshed(ctx, item.ContentID); err != nil { + return "", err + } + return EnrichmentOutcomeNoMatch, nil } + preserveEbookLocalMetadata(item, accumulator) if err := e.persist(ctx, item.ContentID, accumulatedIDs, accumulator); err != nil { - return fmt.Errorf("persisting enrichment for %s: %w", item.ContentID, err) + return "", fmt.Errorf("persisting enrichment for %s: %w", item.ContentID, err) } e.enqueueRemoteArtwork(ctx, item.ContentID, accumulator) e.autoLinkLiteraryWork(ctx, item.ContentID) @@ -369,7 +597,69 @@ func (e *Enricher) enrichWithProviders(ctx context.Context, item enrichmentItemR "people", len(filterEbookPeople(accumulator.People)), ) - return nil + return EnrichmentOutcomeSuccess, nil +} + +func preserveEbookLocalMetadata(item enrichmentItemRow, result *metadata.MetadataResult) { + if result == nil { + return + } + if item.PosterPath != "" && !ebookPosterOwnedByRemoteProvider(item.PosterPath) { + result.PosterPath = "" + result.PosterThumbhash = "" + } + if !strings.EqualFold(strings.TrimSpace(item.Status), "pending") { + return + } + if item.Year > 0 { + result.Year = 0 + } + if item.Overview != "" { + result.Overview = "" + } + if item.ReleaseDate != "" { + result.ReleaseDate = "" + } + if len(item.Genres) > 0 { + result.Genres = nil + } + if len(item.Studios) > 0 { + result.Studios = nil + } + if item.Author != "" { + result.People = nil + } +} + +func ebookPosterOwnedByRemoteProvider(path string) bool { + path = strings.TrimSpace(path) + return isRemoteHTTPImage(path) || strings.HasPrefix(path, ebookMetadataImageProviderID+"/ebooks/") +} + +func classifyEnrichmentError(err error) (EnrichmentErrorClass, time.Duration) { + grpcStatus, ok := status.FromError(err) + if !ok { + return EnrichmentErrorTransient, 0 + } + + switch grpcStatus.Code() { + case codes.ResourceExhausted: + for _, detail := range grpcStatus.Details() { + if retry, ok := detail.(*errdetails.RetryInfo); ok && retry.GetRetryDelay() != nil { + return EnrichmentErrorRateLimited, retry.GetRetryDelay().AsDuration() + } + } + return EnrichmentErrorRateLimited, 0 + case codes.InvalidArgument, + codes.NotFound, + codes.PermissionDenied, + codes.Unauthenticated, + codes.FailedPrecondition, + codes.Unimplemented: + return EnrichmentErrorPermanent, 0 + default: + return EnrichmentErrorTransient, 0 + } } // providerIDOwnerLookup reports the content item (if any) that already owns a @@ -685,46 +975,16 @@ func (e *Enricher) stampLastRefreshed(ctx context.Context, contentID string) err return nil } now := time.Now().UTC() - if _, err := e.pool.Exec(ctx, ` + _, err := e.pool.Exec(ctx, ` UPDATE media_items SET last_refreshed = $1, matched_at = COALESCE(matched_at, $1), status = CASE WHEN status = 'pending' THEN 'matched' ELSE status END WHERE content_id = $2 - `, now, contentID); err != nil { - return err - } - // Success clears the enrichment failure backlog. media_items.refresh_failures - // is intentionally left alone: it belongs to the metadata refresh-debt system. - _, err := e.pool.Exec(ctx, ` - DELETE FROM ebook_enrichment_state WHERE content_id = $1 - `, contentID) + `, now, contentID) return err } -// recordEnrichFailure increments the item's ebook_enrichment_state failure -// counter so claimBatch deprioritizes it on the next sweep and stops claiming -// it at enrichFailureCap. The state is dedicated to ebook enrichment; -// media_items.refresh_failures is owned by the metadata refresh-debt system -// and is never touched here. -func (e *Enricher) recordEnrichFailure(ctx context.Context, item enrichmentItemRow) { - if e == nil || e.pool == nil { - return - } - if _, err := e.pool.Exec(ctx, ` - INSERT INTO ebook_enrichment_state (content_id, failures, updated_at) - VALUES ($1, 1, NOW()) - ON CONFLICT (content_id) DO UPDATE SET - failures = ebook_enrichment_state.failures + 1, - updated_at = NOW() - `, item.ContentID); err != nil { - slog.WarnContext(ctx, "ebook enrichment: failed to record enrichment failure", "component", "ebooks", - "content_id", item.ContentID, - "error", err, - ) - } -} - func (e *Enricher) persistPeople(ctx context.Context, contentID string, people []models.ItemPerson) error { people = filterEbookPeople(people) if len(people) == 0 { diff --git a/internal/ebooks/enrichment_queue.go b/internal/ebooks/enrichment_queue.go new file mode 100644 index 00000000..d25624b3 --- /dev/null +++ b/internal/ebooks/enrichment_queue.go @@ -0,0 +1,387 @@ +package ebooks + +import ( + "context" + "errors" + "fmt" + "strings" + "sync" + "time" + + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/Silo-Server/silo-server/internal/catalog" +) + +const ( + defaultEnrichmentLease = 10 * time.Minute + maxEnrichmentRetry = 24 * time.Hour + transientRetryBase = 5 * time.Minute + skippedRetryHorizon = 15 * time.Minute +) + +type EnrichmentOutcome string + +const ( + EnrichmentOutcomeSuccess EnrichmentOutcome = "success" + EnrichmentOutcomeNoMatch EnrichmentOutcome = "no_match" + EnrichmentOutcomeSkipped EnrichmentOutcome = "skipped" +) + +type EnrichmentErrorClass string + +const ( + EnrichmentErrorTransient EnrichmentErrorClass = "transient" + EnrichmentErrorRateLimited EnrichmentErrorClass = "rate_limited" + EnrichmentErrorPermanent EnrichmentErrorClass = "permanent" +) + +var ErrEnrichmentLeaseLost = errors.New("ebook enrichment lease lost") + +type EnrichmentJob struct { + ContentID string + Attempts int + LastAttemptAt time.Time +} + +type EnrichmentQueue struct { + pool *pgxpool.Pool + + claimsMu sync.Mutex + claims map[string]EnrichmentJob +} + +func NewEnrichmentQueue(pool *pgxpool.Pool) *EnrichmentQueue { + return &EnrichmentQueue{pool: pool} +} + +var enqueueEnrichmentJobQuery = ` + INSERT INTO ebook_enrichment_state ( + content_id, status, priority, next_attempt_at, updated_at + ) + VALUES ($1, 'pending', $2, now(), now()) + ON CONFLICT (content_id) DO UPDATE SET + priority = GREATEST(ebook_enrichment_state.priority, EXCLUDED.priority), + updated_at = now() +` + +func (q *EnrichmentQueue) Enqueue(ctx context.Context, contentID string, priority int) error { + if q == nil || q.pool == nil { + return errors.New("ebook enrichment queue is not configured") + } + if strings.TrimSpace(contentID) == "" { + return errors.New("ebook enrichment content id is required") + } + _, err := q.pool.Exec(ctx, enqueueEnrichmentJobQuery, contentID, priority) + return err +} + +var materializeEnrichmentJobsQuery = ` + INSERT INTO ebook_enrichment_state ( + content_id, status, priority, next_attempt_at, updated_at + ) + SELECT mi.content_id, 'pending', 100, now(), now() + FROM media_items mi + WHERE mi.type = 'ebook' + AND ` + catalog.MangaChapterExclusionWhere("mi") + ` + AND mi.last_refreshed IS NULL + ON CONFLICT (content_id) DO NOTHING +` + +func (q *EnrichmentQueue) MaterializeCandidates(ctx context.Context) error { + if q == nil || q.pool == nil { + return errors.New("ebook enrichment queue is not configured") + } + _, err := q.pool.Exec(ctx, materializeEnrichmentJobsQuery) + return err +} + +var claimEnrichmentJobsQuery = ` + WITH candidates AS ( + SELECT content_id + FROM ebook_enrichment_state + WHERE next_attempt_at <= now() + AND (status = 'pending' OR (status = 'running' AND lease_until < now())) + ORDER BY priority DESC, next_attempt_at, updated_at + FOR UPDATE SKIP LOCKED + LIMIT $1 + ) + UPDATE ebook_enrichment_state state + SET status = 'running', + lease_until = now() + $2::interval, + last_attempt_at = now(), + attempts = attempts + 1, + updated_at = now() + FROM candidates + WHERE state.content_id = candidates.content_id + RETURNING state.content_id, state.attempts, state.last_attempt_at +` + +func (q *EnrichmentQueue) ClaimBatch(ctx context.Context, limit int, leaseDuration time.Duration) ([]EnrichmentJob, error) { + if q == nil || q.pool == nil { + return nil, errors.New("ebook enrichment queue is not configured") + } + if limit <= 0 { + return nil, nil + } + if leaseDuration <= 0 { + leaseDuration = defaultEnrichmentLease + } + + rows, err := q.pool.Query(ctx, claimEnrichmentJobsQuery, limit, postgresInterval(leaseDuration)) + if err != nil { + return nil, err + } + defer rows.Close() + + jobs := make([]EnrichmentJob, 0, limit) + for rows.Next() { + var job EnrichmentJob + if err := rows.Scan(&job.ContentID, &job.Attempts, &job.LastAttemptAt); err != nil { + return nil, err + } + jobs = append(jobs, job) + q.rememberClaim(job) + } + if err := rows.Err(); err != nil { + return nil, err + } + return jobs, nil +} + +var completeEnrichmentJobQuery = ` + UPDATE ebook_enrichment_state + SET status = 'pending', + lease_until = NULL, + completed_at = now(), + next_attempt_at = now() + $3::interval, + outcome = $2, + attempts = 0, + priority = GREATEST(priority, 0), + last_error_class = NULL, + last_error = NULL, + updated_at = now() + WHERE content_id = $1 + AND status = 'running' + AND last_attempt_at = $4 +` + +func (q *EnrichmentQueue) Complete( + ctx context.Context, + contentID string, + outcome EnrichmentOutcome, + refreshAfter time.Duration, +) error { + if q == nil || q.pool == nil { + return errors.New("ebook enrichment queue is not configured") + } + job, ok := q.claimedJob(contentID) + if !ok { + return ErrEnrichmentLeaseLost + } + return q.CompleteClaim(ctx, job, outcome, refreshAfter) +} + +func (q *EnrichmentQueue) CompleteClaim( + ctx context.Context, + job EnrichmentJob, + outcome EnrichmentOutcome, + refreshAfter time.Duration, +) error { + if q == nil || q.pool == nil { + return errors.New("ebook enrichment queue is not configured") + } + if refreshAfter <= 0 { + refreshAfter = enrichmentRefreshHorizon(outcome) + } + if refreshAfter <= 0 { + return fmt.Errorf("unsupported ebook enrichment outcome %q", outcome) + } + + tag, err := q.pool.Exec( + ctx, + completeEnrichmentJobQuery, + job.ContentID, + string(outcome), + postgresInterval(refreshAfter), + job.LastAttemptAt, + ) + if err != nil { + return err + } + if tag.RowsAffected() == 0 { + q.forgetClaim(job) + return ErrEnrichmentLeaseLost + } + q.forgetClaim(job) + return nil +} + +var failEnrichmentJobQuery = ` + UPDATE ebook_enrichment_state + SET status = 'pending', + lease_until = NULL, + next_attempt_at = now() + $4::interval, + outcome = 'failed', + last_error_class = $2, + last_error = $3, + updated_at = now() + WHERE content_id = $1 + AND status = 'running' + AND last_attempt_at = $5 +` + +func (q *EnrichmentQueue) Fail( + ctx context.Context, + contentID string, + errorClass EnrichmentErrorClass, + message string, + retryAfter time.Duration, +) error { + if q == nil || q.pool == nil { + return errors.New("ebook enrichment queue is not configured") + } + job, ok := q.claimedJob(contentID) + if !ok { + return ErrEnrichmentLeaseLost + } + return q.FailClaim(ctx, job, errorClass, message, retryAfter) +} + +func (q *EnrichmentQueue) FailClaim( + ctx context.Context, + job EnrichmentJob, + errorClass EnrichmentErrorClass, + message string, + retryAfter time.Duration, +) error { + if q == nil || q.pool == nil { + return errors.New("ebook enrichment queue is not configured") + } + delay := enrichmentRetryDelay(errorClass, job.Attempts, retryAfter) + tag, err := q.pool.Exec( + ctx, + failEnrichmentJobQuery, + job.ContentID, + string(errorClass), + message, + postgresInterval(delay), + job.LastAttemptAt, + ) + if err != nil { + return err + } + if tag.RowsAffected() == 0 { + q.forgetClaim(job) + return ErrEnrichmentLeaseLost + } + q.forgetClaim(job) + return nil +} + +var releaseEnrichmentJobQuery = ` + UPDATE ebook_enrichment_state + SET status = 'pending', + lease_until = NULL, + attempts = GREATEST(attempts - 1, 0), + updated_at = now() + WHERE content_id = $1 + AND status = 'running' + AND last_attempt_at = $2 +` + +func (q *EnrichmentQueue) Release(ctx context.Context, contentID string) error { + if q == nil || q.pool == nil { + return errors.New("ebook enrichment queue is not configured") + } + job, ok := q.claimedJob(contentID) + if !ok { + return ErrEnrichmentLeaseLost + } + return q.ReleaseClaim(ctx, job) +} + +func (q *EnrichmentQueue) ReleaseClaim(ctx context.Context, job EnrichmentJob) error { + if q == nil || q.pool == nil { + return errors.New("ebook enrichment queue is not configured") + } + tag, err := q.pool.Exec(ctx, releaseEnrichmentJobQuery, job.ContentID, job.LastAttemptAt) + if err != nil { + return err + } + q.forgetClaim(job) + if tag.RowsAffected() == 0 { + return ErrEnrichmentLeaseLost + } + return nil +} + +func (q *EnrichmentQueue) rememberClaim(job EnrichmentJob) { + q.claimsMu.Lock() + defer q.claimsMu.Unlock() + if q.claims == nil { + q.claims = make(map[string]EnrichmentJob) + } + q.claims[job.ContentID] = job +} + +func (q *EnrichmentQueue) claimedJob(contentID string) (EnrichmentJob, bool) { + q.claimsMu.Lock() + defer q.claimsMu.Unlock() + job, ok := q.claims[contentID] + return job, ok +} + +func (q *EnrichmentQueue) forgetClaim(job EnrichmentJob) { + q.claimsMu.Lock() + defer q.claimsMu.Unlock() + current, ok := q.claims[job.ContentID] + if ok && current.LastAttemptAt.Equal(job.LastAttemptAt) { + delete(q.claims, job.ContentID) + } +} + +func enrichmentRefreshHorizon(outcome EnrichmentOutcome) time.Duration { + switch outcome { + case EnrichmentOutcomeSuccess: + return 90 * 24 * time.Hour + case EnrichmentOutcomeNoMatch: + return 30 * 24 * time.Hour + case EnrichmentOutcomeSkipped: + return skippedRetryHorizon + default: + return 0 + } +} + +func enrichmentRetryDelay(errorClass EnrichmentErrorClass, attempts int, retryAfter time.Duration) time.Duration { + switch errorClass { + case EnrichmentErrorRateLimited: + if retryAfter <= 0 { + retryAfter = transientRetryDelay(attempts) + } + return min(retryAfter, maxEnrichmentRetry) + case EnrichmentErrorPermanent: + return 30 * 24 * time.Hour + default: + return transientRetryDelay(attempts) + } +} + +func transientRetryDelay(attempts int) time.Duration { + if attempts < 1 { + attempts = 1 + } + delay := transientRetryBase + for i := 1; i < attempts; i++ { + if delay >= maxEnrichmentRetry/2 { + return maxEnrichmentRetry + } + delay *= 2 + } + return min(delay, maxEnrichmentRetry) +} + +func postgresInterval(duration time.Duration) string { + return fmt.Sprintf("%d microseconds", duration.Microseconds()) +} diff --git a/internal/ebooks/enrichment_queue_test.go b/internal/ebooks/enrichment_queue_test.go new file mode 100644 index 00000000..c161c848 --- /dev/null +++ b/internal/ebooks/enrichment_queue_test.go @@ -0,0 +1,186 @@ +package ebooks + +import ( + "errors" + "fmt" + "os" + "strings" + "testing" + "time" + + "google.golang.org/genproto/googleapis/rpc/errdetails" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "google.golang.org/protobuf/types/known/durationpb" +) + +func TestEnrichmentQueueClaimQueryUsesAtomicLeasedClaims(t *testing.T) { + query := strings.Join(strings.Fields(claimEnrichmentJobsQuery), " ") + + for _, fragment := range []string{ + "WITH candidates AS", + "FOR UPDATE SKIP LOCKED", + "status = 'pending' OR (status = 'running' AND lease_until < now())", + "next_attempt_at <= now()", + "ORDER BY priority DESC, next_attempt_at, updated_at", + "SET status = 'running'", + "lease_until = now() + $2::interval", + "attempts = attempts + 1", + "RETURNING state.content_id, state.attempts", + } { + if !strings.Contains(query, fragment) { + t.Fatalf("claim query missing %q:\n%s", fragment, claimEnrichmentJobsQuery) + } + } +} + +func TestEnrichmentQueueMaterializesUnrefreshedEbooksRegardlessOfPoster(t *testing.T) { + query := strings.Join(strings.Fields(materializeEnrichmentJobsQuery), " ") + + if !strings.Contains(query, "SELECT mi.content_id, 'pending', 100, now(), now()") { + t.Fatalf("newly discovered ebooks must enter ahead of legacy backfill:\n%s", materializeEnrichmentJobsQuery) + } + if !strings.Contains(query, "mi.last_refreshed IS NULL") { + t.Fatalf("materialization must use refresh state as its eligibility gate:\n%s", materializeEnrichmentJobsQuery) + } + if strings.Contains(query, "poster_path") { + t.Fatalf("local or embedded covers must not make ebooks ineligible for metadata enrichment:\n%s", materializeEnrichmentJobsQuery) + } +} + +func TestEnrichmentQueueMigrationSeedsLegacyBacklogAtLowPriority(t *testing.T) { + body, err := os.ReadFile("../../migrations/sql/20260719090000_ebook_enrichment_jobs.sql") + if err != nil { + t.Fatalf("read enrichment queue migration: %v", err) + } + migration := strings.Join(strings.Fields(string(body)), " ") + + for _, fragment := range []string{ + "INSERT INTO ebook_enrichment_state", + "SELECT mi.content_id, 'pending', -100, 0, now(), now()", + "mi.type = 'ebook'", + "mi.last_refreshed IS NULL", + "NOT EXISTS", + "manga_chapters", + "ON CONFLICT (content_id) DO NOTHING", + } { + if !strings.Contains(migration, fragment) { + t.Fatalf("legacy backlog migration missing %q:\n%s", fragment, body) + } + } + if strings.Contains(migration, "poster_path") { + t.Fatalf("legacy backlog must include ebooks with local or embedded covers:\n%s", body) + } +} + +func TestEnrichmentQueueTransitionsKeepDurableRowsAndReleaseLeases(t *testing.T) { + complete := strings.Join(strings.Fields(completeEnrichmentJobQuery), " ") + for _, fragment := range []string{ + "UPDATE ebook_enrichment_state", + "status = 'pending'", + "lease_until = NULL", + "completed_at = now()", + "next_attempt_at = now() + $3::interval", + "outcome = $2", + "attempts = 0", + "priority = GREATEST(priority, 0)", + "AND last_attempt_at = $4", + } { + if !strings.Contains(complete, fragment) { + t.Fatalf("complete query missing %q:\n%s", fragment, completeEnrichmentJobQuery) + } + } + if strings.Contains(strings.ToUpper(complete), "DELETE") { + t.Fatalf("completion must retain durable queue state:\n%s", completeEnrichmentJobQuery) + } + + release := strings.Join(strings.Fields(releaseEnrichmentJobQuery), " ") + for _, fragment := range []string{ + "status = 'pending'", + "lease_until = NULL", + "attempts = GREATEST(attempts - 1, 0)", + "WHERE content_id = $1", + "AND status = 'running'", + "AND last_attempt_at = $2", + } { + if !strings.Contains(release, fragment) { + t.Fatalf("release query missing %q:\n%s", fragment, releaseEnrichmentJobQuery) + } + } + + failure := strings.Join(strings.Fields(failEnrichmentJobQuery), " ") + if !strings.Contains(failure, "AND last_attempt_at = $5") { + t.Fatalf("failure transition must reject stale leases:\n%s", failEnrichmentJobQuery) + } +} + +func TestEnrichmentRetryPolicy(t *testing.T) { + t.Run("outcome refresh horizons", func(t *testing.T) { + if got := enrichmentRefreshHorizon(EnrichmentOutcomeSuccess); got != 90*24*time.Hour { + t.Fatalf("success refresh horizon = %s, want 90 days", got) + } + if got := enrichmentRefreshHorizon(EnrichmentOutcomeNoMatch); got != 30*24*time.Hour { + t.Fatalf("no-match refresh horizon = %s, want 30 days", got) + } + }) + + t.Run("transient failures use capped exponential backoff", func(t *testing.T) { + if got := enrichmentRetryDelay(EnrichmentErrorTransient, 1, 0); got != 5*time.Minute { + t.Fatalf("first transient retry = %s, want 5m", got) + } + if got := enrichmentRetryDelay(EnrichmentErrorTransient, 2, 0); got != 10*time.Minute { + t.Fatalf("second transient retry = %s, want 10m", got) + } + if got := enrichmentRetryDelay(EnrichmentErrorTransient, 20, 0); got != 24*time.Hour { + t.Fatalf("capped transient retry = %s, want 24h", got) + } + }) + + t.Run("rate limits honor provider horizon with a 24 hour cap", func(t *testing.T) { + if got := enrichmentRetryDelay(EnrichmentErrorRateLimited, 1, 45*time.Minute); got != 45*time.Minute { + t.Fatalf("rate-limited retry = %s, want 45m", got) + } + if got := enrichmentRetryDelay(EnrichmentErrorRateLimited, 1, 72*time.Hour); got != 24*time.Hour { + t.Fatalf("capped rate-limited retry = %s, want 24h", got) + } + }) + + t.Run("permanent failures refresh after 30 days", func(t *testing.T) { + if got := enrichmentRetryDelay(EnrichmentErrorPermanent, 1, 0); got != 30*24*time.Hour { + t.Fatalf("permanent retry = %s, want 30 days", got) + } + }) +} + +func TestEnrichmentRetryPolicyClassifiesProviderErrors(t *testing.T) { + limited, err := status.New(codes.ResourceExhausted, "provider quota exhausted").WithDetails( + &errdetails.RetryInfo{RetryDelay: durationpb.New(45 * time.Minute)}, + ) + if err != nil { + t.Fatalf("attach retry info: %v", err) + } + providerErrors := errors.Join(fmt.Errorf("openlibrary search: %w", limited.Err())) + errorClass, retryAfter := classifyEnrichmentError(fmt.Errorf("no metadata obtained: %w", providerErrors)) + if errorClass != EnrichmentErrorRateLimited || retryAfter != 45*time.Minute { + t.Fatalf("rate-limit classification = (%q, %s), want (rate_limited, 45m)", errorClass, retryAfter) + } + + for _, code := range []codes.Code{ + codes.InvalidArgument, + codes.NotFound, + codes.PermissionDenied, + codes.Unauthenticated, + codes.FailedPrecondition, + codes.Unimplemented, + } { + errorClass, retryAfter = classifyEnrichmentError(status.Error(code, "deterministic provider failure")) + if errorClass != EnrichmentErrorPermanent || retryAfter != 0 { + t.Fatalf("%s classification = (%q, %s), want (permanent, 0)", code, errorClass, retryAfter) + } + } + + errorClass, retryAfter = classifyEnrichmentError(status.Error(codes.Unavailable, "provider down")) + if errorClass != EnrichmentErrorTransient || retryAfter != 0 { + t.Fatalf("unavailable classification = (%q, %s), want (transient, 0)", errorClass, retryAfter) + } +} diff --git a/internal/ebooks/enrichment_test.go b/internal/ebooks/enrichment_test.go index bae99c7b..ddba0cb0 100644 --- a/internal/ebooks/enrichment_test.go +++ b/internal/ebooks/enrichment_test.go @@ -168,23 +168,153 @@ func TestRunBatchSkipsFailureRecordingOnCancellation(t *testing.T) { } } -func TestClaimBatchQueryAppliesFailureBackoffAndCap(t *testing.T) { - // Pin the starvation guard: failing items are deprioritized, not retried - // at the head of every sweep, and capped items are never re-claimed. - if !strings.Contains(claimBatchQuery, "LEFT JOIN ebook_enrichment_state ees ON ees.content_id = mi.content_id") { - t.Fatalf("claimBatchQuery must read dedicated ebook enrichment failure state:\n%s", claimBatchQuery) +func TestEnrichmentQueriesKeepEbookAndAudiobookMetadataSeparate(t *testing.T) { + if !strings.Contains(materializeEnrichmentJobsQuery, "mi.type = 'ebook'") { + t.Fatalf("materialization query must target ebooks:\n%s", materializeEnrichmentJobsQuery) } - if !strings.Contains(claimBatchQuery, "COALESCE(ees.failures, 0) < $2") { - t.Fatalf("claimBatchQuery must exclude items at the failure cap:\n%s", claimBatchQuery) + if !strings.Contains(materializeEnrichmentJobsQuery, "NOT EXISTS") || + !strings.Contains(materializeEnrichmentJobsQuery, "manga_chapters") || + !strings.Contains(materializeEnrichmentJobsQuery, "chapter_content_id") { + t.Fatalf("materialization query must exclude manga chapters:\n%s", materializeEnrichmentJobsQuery) } - if !strings.Contains(claimBatchQuery, "ORDER BY COALESCE(ees.failures, 0) ASC, mi.created_at ASC") { - t.Fatalf("claimBatchQuery must claim least-failed items first:\n%s", claimBatchQuery) + if strings.Contains(loadEnrichmentItemsQuery, "narrator") || + strings.Contains(loadEnrichmentItemsQuery, "asin") { + t.Fatalf("ebook load query must not introduce audiobook-only fields:\n%s", loadEnrichmentItemsQuery) } - if strings.Contains(claimBatchQuery, "refresh_failures") { - t.Fatalf("claimBatchQuery must not read media_items.refresh_failures (owned by the metadata refresh-debt system):\n%s", claimBatchQuery) + if !strings.Contains(loadEnrichmentItemsQuery, "ip.kind = 7") { + t.Fatalf("ebook load query must load author credits only:\n%s", loadEnrichmentItemsQuery) } - if enrichFailureCap < 1 { - t.Fatalf("enrichFailureCap = %d, want >= 1", enrichFailureCap) +} + +func TestEnricherRunTransitionsClaimedJobsByOutcome(t *testing.T) { + queue := &fakeEnrichmentQueue{ + jobs: []EnrichmentJob{ + {ContentID: "success", Attempts: 1}, + {ContentID: "no-match", Attempts: 1}, + {ContentID: "skipped", Attempts: 1}, + {ContentID: "failed", Attempts: 3}, + }, + } + items := []enrichmentItemRow{ + {ContentID: "success"}, + {ContentID: "no-match"}, + {ContentID: "skipped"}, + {ContentID: "failed"}, + } + e := &Enricher{ + queue: queue, + batchSize: len(items), + workers: 2, + loadClaimedItemsFn: func(context.Context, []EnrichmentJob) ([]enrichmentItemRow, error) { + return items, nil + }, + enrichClaimedItemFn: func(_ context.Context, item enrichmentItemRow) (EnrichmentOutcome, error) { + switch item.ContentID { + case "success": + return EnrichmentOutcomeSuccess, nil + case "no-match": + return EnrichmentOutcomeNoMatch, nil + case "skipped": + return EnrichmentOutcomeSkipped, nil + default: + return "", errors.New("provider unavailable") + } + }, + } + + enriched, err := e.Run(context.Background()) + if err != nil { + t.Fatalf("Run() error = %v", err) + } + if enriched != 1 { + t.Fatalf("Run() enriched = %d, want 1", enriched) + } + if queue.materializeCalls != 1 || queue.claimCalls != 1 { + t.Fatalf("queue calls: materialize=%d claim=%d, want 1 each", queue.materializeCalls, queue.claimCalls) + } + if queue.claimLimit != len(items) || queue.leaseDuration <= 0 { + t.Fatalf("claim args: limit=%d lease=%s", queue.claimLimit, queue.leaseDuration) + } + for contentID, want := range map[string]EnrichmentOutcome{ + "success": EnrichmentOutcomeSuccess, + "no-match": EnrichmentOutcomeNoMatch, + "skipped": EnrichmentOutcomeSkipped, + } { + if got := queue.completed[contentID]; got != want { + t.Fatalf("completed[%q] = %q, want %q", contentID, got, want) + } + } + if got := queue.failed["failed"]; got != EnrichmentErrorTransient { + t.Fatalf("failed class = %q, want transient", got) + } +} + +func TestEnricherRunReleasesEveryLeaseOnCancellation(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + queue := &fakeEnrichmentQueue{ + jobs: []EnrichmentJob{ + {ContentID: "first", Attempts: 1}, + {ContentID: "second", Attempts: 1}, + }, + } + items := []enrichmentItemRow{{ContentID: "first"}, {ContentID: "second"}} + e := &Enricher{ + queue: queue, + batchSize: len(items), + workers: 1, + loadClaimedItemsFn: func(context.Context, []EnrichmentJob) ([]enrichmentItemRow, error) { + return items, nil + }, + enrichClaimedItemFn: func(context.Context, enrichmentItemRow) (EnrichmentOutcome, error) { + cancel() + return "", context.Canceled + }, + } + + _, err := e.Run(ctx) + if !errors.Is(err, context.Canceled) { + t.Fatalf("Run() error = %v, want context.Canceled", err) + } + sort.Strings(queue.released) + if got := strings.Join(queue.released, ","); got != "first,second" { + t.Fatalf("released leases = %q, want first,second", got) + } + if queue.releaseSawCanceledContext { + t.Fatal("lease release used the canceled run context") + } + if len(queue.failed) != 0 { + t.Fatalf("cancellation recorded item failures: %v", queue.failed) + } +} + +func TestEnricherRunReleasesLeaseWhenCompletionLosesContext(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + queue := &fakeEnrichmentQueue{ + jobs: []EnrichmentJob{{ContentID: "success", Attempts: 1}}, + completeCancel: cancel, + completeErr: context.Canceled, + } + e := &Enricher{ + queue: queue, + batchSize: 1, + workers: 1, + loadClaimedItemsFn: func(context.Context, []EnrichmentJob) ([]enrichmentItemRow, error) { + return []enrichmentItemRow{{ContentID: "success"}}, nil + }, + enrichClaimedItemFn: func(context.Context, enrichmentItemRow) (EnrichmentOutcome, error) { + return EnrichmentOutcomeSuccess, nil + }, + } + + _, err := e.Run(ctx) + if !errors.Is(err, context.Canceled) { + t.Fatalf("Run() error = %v, want context.Canceled", err) + } + if got := strings.Join(queue.released, ","); got != "success" { + t.Fatalf("released leases = %q, want success", got) + } + if queue.releaseSawCanceledContext { + t.Fatal("lease release used the canceled transition context") } } @@ -588,6 +718,109 @@ func TestBuildEbookMetadataRequestCarriesAccumulatedISBN(t *testing.T) { } } +func TestPreservePendingEbookLocalMetadata(t *testing.T) { + item := enrichmentItemRow{ + Status: "pending", + Year: 2021, + Overview: "Embedded description", + ReleaseDate: "2021-03-04", + Genres: []string{"Local genre"}, + Studios: []string{"Local publisher"}, + PosterPath: "local/ebooks/book/poster/original.webp", + Author: "Embedded Author", + } + result := &metadata.MetadataResult{ + HasMetadata: true, + Year: 2022, + Overview: "Remote description", + ReleaseDate: "2022-04-05", + Genres: []string{"Remote genre"}, + Studios: []string{"Remote publisher"}, + PosterPath: "https://example.test/remote.jpg", + Tagline: "Remote-only field", + ProviderIDs: map[string]string{"isbn": "9780306406157"}, + People: []models.ItemPerson{ + {Person: models.Person{Name: "Remote Author"}, Kind: models.PersonKindAuthor}, + }, + } + + preserveEbookLocalMetadata(item, result) + + if result.Year != 0 || result.Overview != "" || result.ReleaseDate != "" || + len(result.Genres) != 0 || len(result.Studios) != 0 || result.PosterPath != "" || + len(result.People) != 0 { + t.Fatalf("remote fields would overwrite scanner metadata: %+v", result) + } + if result.Tagline != "Remote-only field" { + t.Fatalf("empty local field was not enrichable: tagline=%q", result.Tagline) + } + if result.ProviderIDs["isbn"] != "9780306406157" { + t.Fatalf("provider identity was discarded: %v", result.ProviderIDs) + } +} + +func TestPreservePendingEbookLocalMetadataAllowsRefreshReplacement(t *testing.T) { + result := &metadata.MetadataResult{ + HasMetadata: true, + Year: 2022, + Overview: "Corrected remote description", + } + + preserveEbookLocalMetadata(enrichmentItemRow{ + Status: "matched", + Year: 2021, + Overview: "Old remote description", + }, result) + + if result.Year != 2022 || result.Overview != "Corrected remote description" { + t.Fatalf("controlled refresh fields were suppressed: %+v", result) + } +} + +func TestPreserveEbookLocalPosterDuringControlledRefresh(t *testing.T) { + result := &metadata.MetadataResult{ + HasMetadata: true, + PosterPath: "https://example.test/replacement.jpg", + PosterThumbhash: "remote-thumb", + Overview: "Corrected remote description", + } + + preserveEbookLocalMetadata(enrichmentItemRow{ + Status: "matched", + PosterPath: "local/ebooks/book/poster/original.webp", + }, result) + + if result.PosterPath != "" || result.PosterThumbhash != "" { + t.Fatalf("controlled refresh would replace local poster: %+v", result) + } + if result.Overview != "Corrected remote description" { + t.Fatalf("non-artwork refresh field was suppressed: %+v", result) + } +} + +func TestPreserveEbookLocalPosterAllowsProviderOwnedReplacement(t *testing.T) { + for _, current := range []string{ + "", + "https://provider.test/old.jpg", + "ebook-metadata/ebooks/book/poster/original.webp", + } { + result := &metadata.MetadataResult{ + HasMetadata: true, + PosterPath: "https://provider.test/new.jpg", + PosterThumbhash: "new-thumb", + } + + preserveEbookLocalMetadata(enrichmentItemRow{ + Status: "matched", + PosterPath: current, + }, result) + + if result.PosterPath == "" || result.PosterThumbhash == "" { + t.Fatalf("provider-owned poster %q was not replaceable: %+v", current, result) + } + } +} + type fakeEbookMetadataProvider struct { slug string searchErr error @@ -629,6 +862,68 @@ func (f *fakeEbookImageCacher) CacheImage(_ context.Context, req metadata.CacheI }, nil } +type fakeEnrichmentQueue struct { + mu sync.Mutex + jobs []EnrichmentJob + materializeCalls int + claimCalls int + claimLimit int + leaseDuration time.Duration + completed map[string]EnrichmentOutcome + failed map[string]EnrichmentErrorClass + released []string + releaseSawCanceledContext bool + completeCancel context.CancelFunc + completeErr error +} + +func (f *fakeEnrichmentQueue) MaterializeCandidates(context.Context) error { + f.mu.Lock() + defer f.mu.Unlock() + f.materializeCalls++ + return nil +} + +func (f *fakeEnrichmentQueue) ClaimBatch(_ context.Context, limit int, leaseDuration time.Duration) ([]EnrichmentJob, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.claimCalls++ + f.claimLimit = limit + f.leaseDuration = leaseDuration + return append([]EnrichmentJob(nil), f.jobs...), nil +} + +func (f *fakeEnrichmentQueue) Complete(_ context.Context, contentID string, outcome EnrichmentOutcome, _ time.Duration) error { + f.mu.Lock() + defer f.mu.Unlock() + if f.completed == nil { + f.completed = make(map[string]EnrichmentOutcome) + } + f.completed[contentID] = outcome + if f.completeCancel != nil { + f.completeCancel() + } + return f.completeErr +} + +func (f *fakeEnrichmentQueue) Fail(_ context.Context, contentID string, errorClass EnrichmentErrorClass, _ string, _ time.Duration) error { + f.mu.Lock() + defer f.mu.Unlock() + if f.failed == nil { + f.failed = make(map[string]EnrichmentErrorClass) + } + f.failed[contentID] = errorClass + return nil +} + +func (f *fakeEnrichmentQueue) Release(ctx context.Context, contentID string) error { + f.mu.Lock() + defer f.mu.Unlock() + f.releaseSawCanceledContext = f.releaseSawCanceledContext || ctx.Err() != nil + f.released = append(f.released, contentID) + return nil +} + func TestCleanEbookSearchTitle(t *testing.T) { cases := []struct { title, author, want string diff --git a/migrations/sql/20260719090000_ebook_enrichment_jobs.sql b/migrations/sql/20260719090000_ebook_enrichment_jobs.sql new file mode 100644 index 00000000..a749232a --- /dev/null +++ b/migrations/sql/20260719090000_ebook_enrichment_jobs.sql @@ -0,0 +1,70 @@ +-- +goose Up +ALTER TABLE ebook_enrichment_state + ADD COLUMN status text NOT NULL DEFAULT 'pending', + ADD COLUMN priority integer NOT NULL DEFAULT 0, + ADD COLUMN attempts integer NOT NULL DEFAULT 0, + ADD COLUMN next_attempt_at timestamptz NOT NULL DEFAULT now(), + ADD COLUMN lease_until timestamptz, + ADD COLUMN last_attempt_at timestamptz, + ADD COLUMN completed_at timestamptz, + ADD COLUMN outcome text, + ADD COLUMN last_error_class text, + ADD COLUMN last_error text, + ADD CONSTRAINT ebook_enrichment_state_status_check + CHECK (status IN ('pending', 'running')), + ADD CONSTRAINT ebook_enrichment_state_attempts_check + CHECK (attempts >= 0); + +-- Rows already tracked by the former failure counter belong to the legacy +-- backlog. Preserve their history but keep them behind incremental work. +UPDATE ebook_enrichment_state +SET attempts = failures, + priority = -100; + +-- Snapshot the pre-migration library as durable low-priority backfill. Runtime +-- materialization assigns new discoveries priority 100, so this backlog cannot +-- occupy claim slots ahead of incremental or due refresh work. +INSERT INTO ebook_enrichment_state ( + content_id, + status, + priority, + attempts, + next_attempt_at, + updated_at +) +SELECT mi.content_id, 'pending', -100, 0, now(), now() +FROM media_items mi +WHERE mi.type = 'ebook' + AND NOT EXISTS ( + SELECT 1 + FROM manga_chapters mc + WHERE mc.chapter_content_id = mi.content_id + ) + AND mi.last_refreshed IS NULL +ON CONFLICT (content_id) DO NOTHING; + +CREATE INDEX ebook_enrichment_state_claim_idx + ON ebook_enrichment_state (priority DESC, next_attempt_at, updated_at) + WHERE status IN ('pending', 'running'); + +-- +goose Down +DROP INDEX IF EXISTS ebook_enrichment_state_claim_idx; + +-- The former state table only retained positive failure rows. Remove queue rows +-- created by this migration/runtime before restoring that representation. +DELETE FROM ebook_enrichment_state +WHERE failures = 0; + +ALTER TABLE ebook_enrichment_state + DROP CONSTRAINT IF EXISTS ebook_enrichment_state_attempts_check, + DROP CONSTRAINT IF EXISTS ebook_enrichment_state_status_check, + DROP COLUMN IF EXISTS last_error, + DROP COLUMN IF EXISTS last_error_class, + DROP COLUMN IF EXISTS outcome, + DROP COLUMN IF EXISTS completed_at, + DROP COLUMN IF EXISTS last_attempt_at, + DROP COLUMN IF EXISTS lease_until, + DROP COLUMN IF EXISTS next_attempt_at, + DROP COLUMN IF EXISTS attempts, + DROP COLUMN IF EXISTS priority, + DROP COLUMN IF EXISTS status;