feat(artwork): add repository queries to enqueue by current source

This commit is contained in:
Deluan 2026-08-14 12:49:53 -04:00
commit 3eeb453aba
4 changed files with 163 additions and 4 deletions

View file

@ -130,6 +130,12 @@ type ArtworkQueueRepository interface {
EnqueueAllMissing(kind Kind, priority int) (int64, error)
// EnqueueIfMissing inserts only for items with no item_artwork row yet.
EnqueueIfMissing(items ...ArtworkQueueItem) error
// CountBySource reports how many items of a kind currently resolve from the given sources.
// An empty sources slice means every source; "" matches absent state.
CountBySource(kind Kind, sources []string) (int64, error)
// EnqueueBySource inserts queue rows for items of a kind whose current source matches.
// It does not clear existing artwork state: the current image stays until it is replaced.
EnqueueBySource(kind Kind, sources []string, priority int) (int64, error)
// DequeueBatch returns up to n items with retry_at <= now, priority desc, enqueued_at asc.
// Restricted to the given kinds when any are passed, so one kind cannot block another's drain.
DequeueBatch(n int, kinds ...string) ([]ArtworkQueueItem, error)

View file

@ -87,12 +87,39 @@ func (r *artworkQueueRepository) EnqueueIfMissing(items ...model.ArtworkQueueIte
return nil
}
// insertIfNotQueued inserts the rows selected by the given SQL, optionally prefixed by a CTE. DO NOTHING is
// deliberate: a recheck must not bump the priority or retry_at of an already-queued item.
// 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+
` ON CONFLICT (item_kind, item_id, image_type) DO NOTHING`, args...))
` (`+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
}
// 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 {

View file

@ -265,6 +265,83 @@ var _ = Describe("ArtworkQueueRepository", func() {
Expect(got[0].Priority).To(Equal(model.ArtworkPriorityBump), "the existing priority must survive")
})
Describe("EnqueueBySource", func() {
BeforeEach(func() {
artRepo := NewArtworkRepository(context.Background(), GetDBXBuilder())
for _, ia := range []model.ItemArtwork{
{ItemKind: "ar", ItemID: "ar1", ImageType: model.ImageTypePrimary, Hash: "h1", Source: "external:deezer"},
{ItemKind: "ar", ItemID: "ar2", ImageType: model.ImageTypePrimary, Hash: "h2", Source: "external:lastfm"},
{ItemKind: "ar", ItemID: "ar3", ImageType: model.ImageTypePrimary, Hash: "", Source: ""},
{ItemKind: "al", ItemID: "al1", ImageType: model.ImageTypePrimary, Hash: "h4", Source: "external:deezer"},
} {
Expect(artRepo.PutItemArtwork(&ia)).To(Succeed())
}
})
It("enqueues only the matching source within the kind", func() {
n, err := repo.EnqueueBySource(model.KindArtistArtwork, []string{"external:deezer"}, model.ArtworkPriorityRecheck)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(int64(1)), "al1 is a different kind and must not be touched")
got, err := repo.DequeueBatch(10)
Expect(err).ToNot(HaveOccurred())
Expect(slice.Map(got, func(it model.ArtworkQueueItem) string { return it.ItemID })).To(ConsistOf("ar1"))
})
It("treats the empty source as absent", func() {
n, err := repo.EnqueueBySource(model.KindArtistArtwork, []string{""}, model.ArtworkPriorityRecheck)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(int64(1)))
got, _ := repo.DequeueBatch(10)
Expect(slice.Map(got, func(it model.ArtworkQueueItem) string { return it.ItemID })).To(ConsistOf("ar3"))
})
It("enqueues every source when none is given", func() {
n, err := repo.EnqueueBySource(model.KindArtistArtwork, nil, model.ArtworkPriorityRecheck)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(int64(3)))
})
It("leaves the current artwork state in place", func() {
_, err := repo.EnqueueBySource(model.KindArtistArtwork, []string{"external:deezer"}, model.ArtworkPriorityRecheck)
Expect(err).ToNot(HaveOccurred())
artRepo := NewArtworkRepository(context.Background(), GetDBXBuilder())
ia, err := artRepo.GetItemArtwork(model.KindArtistArtwork, "ar1", model.ImageTypePrimary)
Expect(err).ToNot(HaveOccurred())
Expect(ia.Hash).To(Equal("h1"), "the current image must survive until it is replaced")
Expect(ia.Source).To(Equal("external:deezer"))
})
It("does not disturb an already-queued row", func() {
Expect(repo.Enqueue(item("ar", "ar1", model.ArtworkPriorityBump))).To(Succeed())
n, err := repo.EnqueueBySource(model.KindArtistArtwork, []string{"external:deezer"}, model.ArtworkPriorityRecheck)
Expect(err).ToNot(HaveOccurred())
Expect(n).To(BeZero())
got, _ := repo.DequeueBatch(10)
Expect(got).To(HaveLen(1))
Expect(got[0].Priority).To(Equal(model.ArtworkPriorityBump))
})
It("counts without enqueueing", func() {
n, err := repo.CountBySource(model.KindArtistArtwork, []string{"external:deezer"})
Expect(err).ToNot(HaveOccurred())
Expect(n).To(Equal(int64(1)))
queued, err := repo.Count()
Expect(err).ToNot(HaveOccurred())
Expect(queued).To(BeZero(), "CountBySource must not enqueue")
})
It("counts the absent source and every source", func() {
Expect(repo.CountBySource(model.KindArtistArtwork, []string{""})).To(Equal(int64(1)))
Expect(repo.CountBySource(model.KindArtistArtwork, nil)).To(Equal(int64(3)))
})
})
It("does not disturb an already-queued entity when enqueueing missing rows", func() {
Expect(repo.Enqueue(item("al", albumRadioactivity.ID, model.ArtworkPriorityBump))).To(Succeed())

View file

@ -216,6 +216,55 @@ func (m *MockArtworkQueueRepo) EnqueueStaleAbsent(kind model.Kind, attemptedBefo
return inserted, nil
}
// matchingSource mirrors the SQL filter: no sources means every source, "" the absent state.
func (m *MockArtworkQueueRepo) matchingSource(kind model.Kind, sources []string) []model.ItemArtwork {
if m.ItemArtworkSource == nil {
return nil
}
var res []model.ItemArtwork
for _, ia := range m.ItemArtworkSource.ItemData {
if ia.ItemKind == kind.Prefix() && (len(sources) == 0 || slices.Contains(sources, ia.Source)) {
res = append(res, ia)
}
}
return res
}
func (m *MockArtworkQueueRepo) CountBySource(kind model.Kind, sources []string) (int64, error) {
m.mu.Lock()
defer m.mu.Unlock()
if m.Err != nil {
return 0, m.Err
}
return int64(len(m.matchingSource(kind, sources))), nil
}
func (m *MockArtworkQueueRepo) EnqueueBySource(kind model.Kind, sources []string, priority int) (int64, error) {
m.mu.Lock()
defer m.mu.Unlock()
if m.Err != nil {
return 0, m.Err
}
now := time.Now()
var inserted int64
for _, ia := range m.matchingSource(kind, sources) {
k := iaKey(ia.ItemKind, ia.ItemID, ia.ImageType)
if _, ok := m.Data[k]; ok { // DO NOTHING: never touch existing queue rows
continue
}
m.Data[k] = model.ArtworkQueueItem{
ItemKind: ia.ItemKind,
ItemID: ia.ItemID,
ImageType: ia.ImageType,
Priority: priority,
RetryAt: now,
EnqueuedAt: now,
}
inserted++
}
return inserted, nil
}
// EnqueueMissing mirrors the SQL set-difference insert: ExistingIDs[kind] minus ItemArtworkSource.
func (m *MockArtworkQueueRepo) EnqueueAllMissing(kind model.Kind, priority int) (int64, error) {
m.mu.Lock()