diff --git a/internal/ebooks/enrichment_queue.go b/internal/ebooks/enrichment_queue.go index 024b7e48..68f7bcf0 100644 --- a/internal/ebooks/enrichment_queue.go +++ b/internal/ebooks/enrichment_queue.go @@ -129,6 +129,9 @@ const mergeEbookProtectedFieldsSQL = ` ) ` +// Legacy-lane rows (priority < 0) may only leave the lane through a terminal +// outcome (Complete/Discard); enqueue, fail, and release must preserve the +// lane so the paced backfill task stays the sole drain for the backlog. var enqueueEnrichmentJobQuery = ` INSERT INTO ebook_enrichment_state ( content_id, status, priority, next_attempt_at, protected_fields, updated_at @@ -145,6 +148,7 @@ var enqueueEnrichmentJobQuery = ` END, priority = CASE WHEN ebook_enrichment_state.status = 'running' THEN ebook_enrichment_state.priority + WHEN ebook_enrichment_state.priority < 0 THEN ebook_enrichment_state.priority ELSE GREATEST(ebook_enrichment_state.priority, EXCLUDED.priority) END, next_attempt_at = CASE @@ -670,7 +674,7 @@ var failEnrichmentJobQuery = ` WHEN requeue_requested THEN now() ELSE now() + $4::interval END, - priority = CASE WHEN requeue_requested THEN 100 ELSE priority END, + priority = CASE WHEN requeue_requested AND priority >= 0 THEN 100 ELSE priority END, requeue_requested = false, outcome = 'failed', last_error_class = $2, @@ -720,7 +724,7 @@ var releaseEnrichmentJobQuery = ` claim_token = NULL, attempts = GREATEST(attempts - 1, 0), next_attempt_at = CASE WHEN requeue_requested THEN now() ELSE next_attempt_at END, - priority = CASE WHEN requeue_requested THEN 100 ELSE priority END, + priority = CASE WHEN requeue_requested AND priority >= 0 THEN 100 ELSE priority END, requeue_requested = false, updated_at = now() WHERE content_id = $1 diff --git a/internal/ebooks/enrichment_queue_test.go b/internal/ebooks/enrichment_queue_test.go index ede1d664..58e82695 100644 --- a/internal/ebooks/enrichment_queue_test.go +++ b/internal/ebooks/enrichment_queue_test.go @@ -317,6 +317,7 @@ func TestEnrichmentQueueEnqueueMakesPendingRowsDueAndDefersRunningRequeue(t *tes "ON CONFLICT (content_id) DO UPDATE SET", "WHEN ebook_enrichment_state.status = 'running' THEN ebook_enrichment_state.next_attempt_at", "ELSE now()", + "WHEN ebook_enrichment_state.priority < 0 THEN ebook_enrichment_state.priority", "requeue_requested = ebook_enrichment_state.requeue_requested OR ebook_enrichment_state.status = 'running'", } { if !strings.Contains(query, fragment) { @@ -328,6 +329,28 @@ func TestEnrichmentQueueEnqueueMakesPendingRowsDueAndDefersRunningRequeue(t *tes } } +func TestEnrichmentQueueKeepsLegacyLaneRowsUntilTerminalOutcome(t *testing.T) { + enqueue := strings.Join(strings.Fields(enqueueEnrichmentJobQuery), " ") + if !strings.Contains(enqueue, "WHEN ebook_enrichment_state.priority < 0 THEN ebook_enrichment_state.priority") { + t.Fatalf("enqueue must keep legacy-lane rows in the legacy lane:\n%s", enqueueEnrichmentJobQuery) + } + for _, tt := range []struct { + name string + query string + }{ + {name: "fail", query: failEnrichmentJobQuery}, + {name: "release", query: releaseEnrichmentJobQuery}, + } { + query := strings.Join(strings.Fields(tt.query), " ") + if strings.Contains(query, "WHEN requeue_requested THEN 100") { + t.Fatalf("%s requeue must not promote legacy-lane rows without a terminal outcome:\n%s", tt.name, tt.query) + } + if !strings.Contains(query, "WHEN requeue_requested AND priority >= 0 THEN 100 ELSE priority END") { + t.Fatalf("%s requeue must stay within the row's lane:\n%s", tt.name, tt.query) + } + } +} + func TestEnrichmentQueueMigrationIsCrashSafeAndSeedsAllLegacyEbooks(t *testing.T) { body, err := os.ReadFile("../../migrations/sql/20260719090000_ebook_enrichment_jobs.sql") if err != nil { @@ -436,7 +459,7 @@ func TestEnrichmentQueueTransitionsKeepDurableRowsAndReleaseLeases(t *testing.T) failure := strings.Join(strings.Fields(failEnrichmentJobQuery), " ") for _, fragment := range []string{ "WHEN requeue_requested THEN now()", - "WHEN requeue_requested THEN 100 ELSE priority END", + "WHEN requeue_requested AND priority >= 0 THEN 100 ELSE priority END", "AND claim_token = $5", } { if !strings.Contains(failure, fragment) {