navidrome/persistence/artwork_queue_repository.go
Deluan Quintão ffc68e29db
feat(cli): add artwork cancel to call off queued artwork work (#6006)
* feat(cli): add `artwork cancel` to call off queued artwork work

A bulk backfill had no off switch. Changing an artwork setting bumps the config
fingerprint, which enqueues every entity in the library, and the only way to stop
it was to turn agents off -- which changes the fingerprint again and enqueues a
second full backfill. The escape hatch was the trap.

`artwork cancel` deletes pending queue rows selected by --kind and/or --priority,
with the --dry-run/confirm/-y flow `reprocess` already uses. Cancelling by
priority is the point: it drops a runaway backfill while leaving the bump-priority
rows an operator queued by hand.

It only touches the queue. Resolved artwork and the item_artwork state behind
`artwork explain` are left alone, and the trace of why a cancelled item last
failed goes with its row. Preserving that trace would mean writing it to
last_failure, which `explain` prints under "Gave up after" -- reporting a
cancellation as an exhausted retry budget. The help text says the trace is
discarded instead.

Two limits the help text states, because neither is guessable: work already
dequeued is not interrupted, and an item with no artwork state yet can be queued
again by the hourly missing-artwork recheck. Cancel calls off queued work; it
does not stop the worker.

--kind validates against RefreshableKinds, not the RecheckKinds `reprocess` uses:
the queue holds media file rows, so --all has to reach them. PurgeQueued follows
the repository's naming rule -- it finds its own rows and reports how many went --
and ignores retry_at, since a row still backing off is pending work. The preview
reuses CountByKindAndPriority rather than adding a counter. reprocessConfirm
became confirmUnlessYes(yes, in, verb) now that two commands prompt.

* refactor(cli): share the artwork queue filter between the preview and the delete

Follow-up cleanup on the previous commit; no change to what the command does,
apart from --all, noted below.

The "which rows does cancel touch" predicate was written three times: once as SQL
in PurgeQueued, once in Go in cmd's matchingQueueStats, and once more in the mock.
The preview and the delete could therefore drift, and the mock would keep the
tests green while they did. persistence now has one artworkQueueFilter, shared by
PurgeQueued and a new CountQueued, and cmd does no filtering at all.

That also makes the preview cheaper. It counted the whole queue and filtered in
Go, so `artwork cancel --kind al` scanned every row of every kind to print a
handful. CountQueued pushes the filter into SQL, which the drain index serves as a
range seek. CountByKindAndPriority is gone: it is CountQueued(nil, nil).

--all now selects with an empty filter instead of enumerating RefreshableKinds.
It is what the flag help already claimed, and the enumeration was narrower than
its own documentation -- a queue row whose item_kind this build does not know
survived `--all` with no flag combination able to remove it. It also restores
SQLite's truncate path: measured with EXPLAIN QUERY PLAN, a bare DELETE plans to
nothing, while `WHERE (1=1)` -- which an empty squirrel And renders -- plans to a
full index scan. A test pins the filter's emptiness so that cannot regress
silently.

Also folded together three copies of the parse-and-dedup loop (parseAll), two
copies of the queue-stats table (printQueueStats, now shared with `artwork
status`), two copies of the stat sum (queueTotal), and four copies of the
kind-to-prefix mapping (model.KindPrefixes). The PurgeQueued specs became one
DescribeTable that asserts count and delete agree on every selection.

* docs(cli): say when `artwork cancel` evaluates its selection

The help text covered the two limits that surprise an operator after the fact, but
not the one that bites during the prompt: the count is a preview, and the filters
run again on confirm. A scan or a manual refresh landing in between is cancelled
without ever appearing in the table the operator agreed to.

Deleting only the previewed rows was considered and rejected. The exposure is one
item re-resolving on next view instead of immediately: clearing an item's artwork
state is what every recovery path selects on, so a lost Bump row from
artwork.Refresh comes back at the same priority via provisional() on the next
request, and otherwise within the hour via EnqueueAllMissing. Buying a guarantee
against that costs the truncate path on --all, the flag that exists for a
29k-item backfill.

* refactor(cli): share one set of flag targets across the artwork subcommands

reprocess and cancel each declared their own kinds/all/dry-run/yes variables, but
cobra only ever parses the one subcommand being run, so the two sets could never
hold values at the same time. backup.go already binds one backupDir across two
subcommands and one force across two more; this follows that.

Ten package-level variables become six. Each command keeps its own help string
and its own valid-kind list, so --kind still reports RecheckKinds for reprocess
and RefreshableKinds for cancel, and --source and --priority stay registered only
on the command that has them.

The priority lookup table is now knownPriorities, freeing the artworkPriorities
name for the flag. The new name also reads better against priorityName's fallback
for a value it does not know.
2026-08-21 14:03:59 -04:00

242 lines
9.3 KiB
Go

package persistence
import (
"cmp"
"context"
"fmt"
"slices"
"strings"
"time"
. "github.com/Masterminds/squirrel"
"github.com/navidrome/navidrome/model"
"github.com/navidrome/navidrome/utils/slice"
"github.com/pocketbase/dbx"
)
// Keeps each multi-row insert under SQLite's bind-variable limit (at most 7 vars per row).
const enqueueChunkSize = 100
// Every insert writes these, in this order; the INSERT..SELECT forms must project them to match.
// DequeueBatch also selects exactly these, to leave the drain's rows free of the trace it never reads.
var enqueueColumns = []string{"item_kind", "item_id", "image_type", "priority", "attempts", "retry_at", "enqueued_at"}
type artworkQueueRepository struct {
sqlRepository
}
func NewArtworkQueueRepository(ctx context.Context, db dbx.Builder) model.ArtworkQueueRepository {
r := &artworkQueueRepository{}
r.ctx = ctx
r.db = db
r.tableName = "artwork_queue"
return r
}
func (r *artworkQueueRepository) Get(kind model.Kind, id, imageType string) (*model.ArtworkQueueItem, error) {
var res model.ArtworkQueueItem
err := r.queryOne(Select("*").From(r.tableName).
Where(Eq{"item_kind": kind.Prefix(), "item_id": id, "image_type": imageType}), &res)
if err != nil {
return nil, err
}
return &res, nil
}
// Enqueue starts a fresh lifecycle: it resets enqueued_at (so a fresh request does not inherit an old
// row's spent retry budget) and clears trace (so explain does not show a prior failure at attempts 0).
func (r *artworkQueueRepository) Enqueue(items ...model.ArtworkQueueItem) error {
return r.enqueue(`ON CONFLICT (item_kind, item_id, image_type) DO UPDATE SET
priority = MAX(priority, excluded.priority), retry_at = excluded.retry_at,
attempts = 0, enqueued_at = excluded.enqueued_at, trace = '[]'`, items)
}
func (r *artworkQueueRepository) EnqueuePreservingBackoff(items ...model.ArtworkQueueItem) error {
return r.enqueue(`ON CONFLICT (item_kind, item_id, image_type) DO UPDATE SET
priority = MAX(priority, excluded.priority)`, items)
}
func (r *artworkQueueRepository) EnqueueStaleAbsent(kind model.Kind, attemptedBefore time.Time) (int64, error) {
now := time.Now()
return r.insertIfNotQueued("", `SELECT item_kind, item_id, image_type, ?, 0, ?, ?
FROM `+itemArtworkTable+` WHERE item_kind = ? AND hash = '' AND attempted_at < ?`,
model.ArtworkPriorityRecheck, now, now, kind.Prefix(), attemptedBefore)
}
func (r *artworkQueueRepository) EnqueueAllMissing(kind model.Kind, priority int) (int64, error) {
entityTable, ok := artworkOwnerTables[kind]
if !ok {
return 0, fmt.Errorf("artwork queue: no entity table for kind %q", kind.Prefix())
}
now := time.Now()
return r.insertIfNotQueued("", `SELECT ?, id, ?, ?, 0, ?, ?
FROM `+entityTable+`
WHERE id NOT IN (SELECT item_id FROM `+itemArtworkTable+` WHERE item_kind = ?)`,
kind.Prefix(), model.ImageTypePrimary, priority, now, now, kind.Prefix())
}
func (r *artworkQueueRepository) EnqueueIfMissing(items ...model.ArtworkQueueItem) error {
now := time.Now()
for chunk := range slices.Chunk(items, enqueueChunkSize) {
rows := make([]string, 0, len(chunk))
args := make([]any, 0, len(chunk)*4+2)
for _, it := range chunk {
rows = append(rows, "(?,?,?,?)")
args = append(args, it.ItemKind, it.ItemID, cmp.Or(it.ImageType, model.ImageTypePrimary), it.Priority)
}
args = append(args, now, now)
_, err := r.insertIfNotQueued(
`WITH new_items(item_kind, item_id, image_type, priority) AS (VALUES `+strings.Join(rows, ",")+`) `,
`SELECT n.item_kind, n.item_id, n.image_type, n.priority, 0, ?, ?
FROM new_items n
WHERE NOT EXISTS (
SELECT 1 FROM `+itemArtworkTable+` ia
WHERE ia.item_kind = n.item_kind AND ia.item_id = n.item_id AND ia.image_type = n.image_type)`,
args...)
if err != nil {
return err
}
}
return nil
}
// DO NOTHING is deliberate: a recheck must not bump the priority or retry_at of an already-queued item.
const skipIfQueued = ` ON CONFLICT (item_kind, item_id, image_type) DO NOTHING`
// insertIfNotQueued inserts the rows selected by the given SQL, optionally prefixed by a CTE.
func (r *artworkQueueRepository) insertIfNotQueued(with, sql string, args ...any) (int64, error) {
return r.executeSQL(Expr(with+`INSERT INTO `+r.tableName+
` (`+strings.Join(enqueueColumns, ", ")+`) `+sql+skipIfQueued, args...))
}
// artworkSourceFilter selects item_artwork rows of a kind; no sources means every source, "" the absent state.
func artworkSourceFilter(kind model.Kind, sources []string) Sqlizer {
f := And{Eq{"item_kind": kind.Prefix()}}
if len(sources) > 0 {
f = append(f, Eq{"source": sources})
}
return f
}
func (r *artworkQueueRepository) CountBySource(kind model.Kind, sources []string) (int64, error) {
var res struct{ Count int64 }
err := r.queryOne(Select("count(*) as count").From(itemArtworkTable).
Where(artworkSourceFilter(kind, sources)), &res)
return res.Count, err
}
func (r *artworkQueueRepository) SourcesInUse(kind model.Kind) ([]string, error) {
var res []struct{ Source string }
err := r.queryAll(Select("distinct source").From(itemArtworkTable).
Where(Eq{"item_kind": kind.Prefix()}), &res)
if err != nil {
return nil, err
}
return slice.Map(res, func(s struct{ Source string }) string { return s.Source }), nil
}
// EnqueueBySource deliberately leaves item_artwork alone: clearing state in bulk would blank the
// library's artwork until every item is resolved again.
func (r *artworkQueueRepository) EnqueueBySource(kind model.Kind, sources []string, priority int) (int64, error) {
now := time.Now()
sel := Select("item_kind", "item_id", "image_type").
Column(Expr("?", priority)).Column("0").Column(Expr("?", now)).Column(Expr("?", now)).
From(itemArtworkTable).Where(artworkSourceFilter(kind, sources))
return r.executeSQL(Insert(r.tableName).Columns(enqueueColumns...).Select(sel).Suffix(skipIfQueued))
}
func (r *artworkQueueRepository) enqueue(conflict string, items []model.ArtworkQueueItem) error {
now := time.Now()
for chunk := range slices.Chunk(items, enqueueChunkSize) {
ins := Insert(r.tableName).Columns(enqueueColumns...)
for _, it := range chunk {
ins = ins.Values(it.ItemKind, it.ItemID, cmp.Or(it.ImageType, model.ImageTypePrimary), it.Priority, 0, now, now)
}
ins = ins.Suffix(conflict)
if _, err := r.executeSQL(ins); err != nil {
return err
}
}
return nil
}
func (r *artworkQueueRepository) DequeueBatch(n int, kinds ...string) ([]model.ArtworkQueueItem, error) {
sel := Select(enqueueColumns...).From(r.tableName).
Where(LtOrEq{"retry_at": time.Now()}).
OrderBy("priority DESC", "enqueued_at ASC").
Limit(uint64(n))
if len(kinds) > 0 {
sel = sel.Where(Eq{"item_kind": kinds})
}
var res []model.ArtworkQueueItem
err := r.queryAll(sel, &res)
return res, err
}
func (r *artworkQueueRepository) MarkFailedIfUnchanged(kind, id, imageType string, seenRetryAt, retryAt time.Time, trace string) error {
upd := Update(r.tableName).
Set("attempts", Expr("attempts + 1")).
Set("retry_at", retryAt).
Set("trace", trace).
Where(Eq{"item_kind": kind, "item_id": id, "image_type": imageType, "retry_at": seenRetryAt})
_, err := r.executeSQL(upd)
return err
}
func (r *artworkQueueRepository) DeleteIfUnchanged(kind, id, imageType string, retryAt time.Time) error {
return r.delete(Eq{"item_kind": kind, "item_id": id, "image_type": imageType, "retry_at": retryAt})
}
func (r *artworkQueueRepository) PurgeDangling() (int64, error) {
return purgeDangling(r.sqlRepository)
}
// artworkQueueFilter returns no conditions for an empty filter, so an unfiltered DELETE keeps
// SQLite's truncate path. It ignores retry_at: a backing-off row is pending work too.
func artworkQueueFilter(kinds []model.Kind, priorities []int) And {
var f And
if len(kinds) > 0 {
f = append(f, Eq{"item_kind": model.KindPrefixes(kinds)})
}
if len(priorities) > 0 {
f = append(f, Eq{"priority": priorities})
}
return f
}
// CountQueued shares its filter with PurgeQueued, so a preview cannot count rows the delete misses.
func (r *artworkQueueRepository) CountQueued(kinds []model.Kind, priorities []int) ([]model.ArtworkQueueStat, error) {
sel := Select("item_kind", "priority", "count(*) as count").From(r.tableName).
GroupBy("item_kind", "priority").OrderBy("item_kind", "priority desc")
if f := artworkQueueFilter(kinds, priorities); len(f) > 0 {
sel = sel.Where(f)
}
var res []model.ArtworkQueueStat
err := r.queryAll(sel, &res)
return res, err
}
func (r *artworkQueueRepository) PurgeQueued(kinds []model.Kind, priorities []int) (int64, error) {
del := Delete(r.tableName)
if f := artworkQueueFilter(kinds, priorities); len(f) > 0 {
del = del.Where(f)
}
return r.executeSQL(del)
}
func (r *artworkQueueRepository) Count() (int64, error) {
var res struct{ Count int64 }
err := r.queryOne(Select("count(*) as count").From(r.tableName), &res)
return res.Count, err
}
// CountAbsent matches EnqueueStaleAbsent on hash, so the stale count is what a recheck would queue.
func (r *artworkQueueRepository) CountAbsent(kind model.Kind, attemptedBefore time.Time) (model.ArtworkAbsentStat, error) {
var res model.ArtworkAbsentStat
err := r.queryOne(Select("count(*) as total").
Column(Expr("coalesce(sum(attempted_at < ?), 0) as stale", attemptedBefore)).
From(itemArtworkTable).Where(Eq{"item_kind": kind.Prefix(), "hash": ""}), &res)
return res, err
}
var _ model.ArtworkQueueRepository = (*artworkQueueRepository)(nil)