From 3eeb453aba43f5c4b5af3f4bb964f57b40b434a1 Mon Sep 17 00:00:00 2001 From: Deluan Date: Fri, 14 Aug 2026 12:49:53 -0400 Subject: [PATCH] feat(artwork): add repository queries to enqueue by current source --- model/artwork.go | 6 ++ persistence/artwork_queue_repository.go | 35 ++++++++- persistence/artwork_queue_repository_test.go | 77 ++++++++++++++++++++ tests/mock_artwork_queue_repo.go | 49 +++++++++++++ 4 files changed, 163 insertions(+), 4 deletions(-) diff --git a/model/artwork.go b/model/artwork.go index 87b424f33..2612f6846 100644 --- a/model/artwork.go +++ b/model/artwork.go @@ -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) diff --git a/persistence/artwork_queue_repository.go b/persistence/artwork_queue_repository.go index e5469fbea..526d9c9d1 100644 --- a/persistence/artwork_queue_repository.go +++ b/persistence/artwork_queue_repository.go @@ -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 { diff --git a/persistence/artwork_queue_repository_test.go b/persistence/artwork_queue_repository_test.go index 6638a204e..816b314a4 100644 --- a/persistence/artwork_queue_repository_test.go +++ b/persistence/artwork_queue_repository_test.go @@ -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()) diff --git a/tests/mock_artwork_queue_repo.go b/tests/mock_artwork_queue_repo.go index 1b097ca32..0eeef465a 100644 --- a/tests/mock_artwork_queue_repo.go +++ b/tests/mock_artwork_queue_repo.go @@ -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()