feat(ebooks): add durable metadata enrichment queue

This commit is contained in:
rxwatcher
2026-07-22 13:43:57 +02:00
parent 1ab85d18ea
commit ef7eedf3fc
5 changed files with 1299 additions and 101 deletions
+348 -88
View File
@@ -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 {
+387
View File
@@ -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())
}
+186
View File
@@ -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)
}
}
+308 -13
View File
@@ -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
@@ -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;