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