release / go (push) Successful in 2m33s
release / web (push) Successful in 1m41s
release / govulncheck (push) Successful in 22s
release / integration (push) Successful in 6m0s
release / android (push) Successful in 6m7s
release / Build signed APK (releases and dev) (push) Successful in 6m8s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 1m43s
release / Verify release artifacts (tag releases only) (push) Skipped
Lidarr's metadata is keyed by MusicBrainz release group, but re-acquisition requested albums by their release id, so every add came back "not found". - Sweeper requests an album by its tag-supplied release group, else the one MusicBrainz names (cached onto the album). An album MusicBrainz cannot name is skipped and counted, with no attempt spent. - Reconciler: an add refused as not found re-reads the request's album id as a release (library first, then MusicBrainz), rewrites the request to the group and adds again. This repairs the requests already stored. - Completion matches an album by release id or release group, and only once a track of it is on disk, so a re-acquisition request no longer completes against the row of the album it is trying to bring back. Closes #5241. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
368 lines
12 KiB
Go
368 lines
12 KiB
Go
package lidarrrequests
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgtype"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
|
|
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
|
|
"git.fabledsword.com/bvandeusen/minstrel/internal/eventbus"
|
|
"git.fabledsword.com/bvandeusen/minstrel/internal/lidarr"
|
|
"git.fabledsword.com/bvandeusen/minstrel/internal/lidarrconfig"
|
|
"git.fabledsword.com/bvandeusen/minstrel/internal/tags"
|
|
)
|
|
|
|
// Reconciler is a background worker that periodically scans approved
|
|
// lidarr_requests. For rows whose Lidarr add hasn't been confirmed it
|
|
// (re)sends the add idempotently until it sticks; for confirmed rows it
|
|
// transitions them to completed when their target track/album/artist has
|
|
// appeared in the local library (matched by MBID). Errors per-row are
|
|
// logged at WARN and do not abort the tick.
|
|
type Reconciler struct {
|
|
pool *pgxpool.Pool
|
|
lidarrCfg *lidarrconfig.Service
|
|
clientFn func() *lidarr.Client
|
|
logger *slog.Logger
|
|
bus *eventbus.Bus
|
|
tick time.Duration
|
|
batch int32
|
|
// releaseGroup names the MusicBrainz release group of a release id, to
|
|
// repair a request stored under a release id Lidarr does not index
|
|
// (#5241). A field so tests can stand in for MusicBrainz.
|
|
releaseGroup func(ctx context.Context, releaseMBID string) (string, error)
|
|
}
|
|
|
|
// NewReconciler constructs a Reconciler with production defaults:
|
|
// 5-minute tick, batch size 50. clientFn is the per-tick Lidarr client
|
|
// factory (same pattern as Service) so config changes apply without a
|
|
// restart; pass nil to disable the add-retry path. Pass nil for bus when
|
|
// no SSE publishing is desired (tests, or environments without the event
|
|
// stream enabled).
|
|
func NewReconciler(pool *pgxpool.Pool, cfg *lidarrconfig.Service, clientFn func() *lidarr.Client, logger *slog.Logger, bus *eventbus.Bus) *Reconciler {
|
|
if clientFn == nil {
|
|
clientFn = func() *lidarr.Client { return nil }
|
|
}
|
|
return &Reconciler{
|
|
pool: pool,
|
|
lidarrCfg: cfg,
|
|
clientFn: clientFn,
|
|
logger: logger,
|
|
bus: bus,
|
|
tick: 5 * time.Minute,
|
|
batch: 50,
|
|
|
|
releaseGroup: tags.ReleaseGroupForRelease,
|
|
}
|
|
}
|
|
|
|
// publishCompleted broadcasts a request.status_changed event scoped to
|
|
// the original requester. No-op when bus is nil.
|
|
func (r *Reconciler) publishCompleted(row dbq.LidarrRequest) {
|
|
if r.bus == nil {
|
|
return
|
|
}
|
|
r.bus.Publish(eventbus.Event{
|
|
Kind: "request.status_changed",
|
|
UserID: formatUUIDForBus(row.UserID),
|
|
Data: map[string]any{
|
|
"request_id": formatUUIDForBus(row.ID),
|
|
"status": "completed",
|
|
"kind": string(row.Kind),
|
|
},
|
|
})
|
|
}
|
|
|
|
// formatUUIDForBus renders a pgtype.UUID as the canonical 8-4-4-4-12 hex
|
|
// string. Mirrors helpers in internal/api/convert.go and other packages;
|
|
// inlined here to avoid a back-edge dependency from lidarrrequests onto
|
|
// the api layer.
|
|
func formatUUIDForBus(u pgtype.UUID) string {
|
|
if !u.Valid {
|
|
return ""
|
|
}
|
|
const hex = "0123456789abcdef"
|
|
out := make([]byte, 36)
|
|
pos := 0
|
|
for i, b := range u.Bytes {
|
|
if i == 4 || i == 6 || i == 8 || i == 10 {
|
|
out[pos] = '-'
|
|
pos++
|
|
}
|
|
out[pos] = hex[b>>4]
|
|
out[pos+1] = hex[b&0xf]
|
|
pos += 2
|
|
}
|
|
return string(out)
|
|
}
|
|
|
|
// Run blocks until ctx is cancelled, ticking every r.tick. Errors from
|
|
// tickOnce are logged at WARN and never propagated.
|
|
func (r *Reconciler) Run(ctx context.Context) {
|
|
t := time.NewTicker(r.tick)
|
|
defer t.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-t.C:
|
|
if err := r.tickOnce(ctx); err != nil {
|
|
r.logger.Warn("lidarrrequests: reconciler tick failed", "err", err)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// tickOnce loads approved requests (oldest first, up to r.batch) and
|
|
// transitions any whose target MBID is now present in the local library.
|
|
// It is exported-for-tests via the lowercase name (package-internal).
|
|
func (r *Reconciler) tickOnce(ctx context.Context) error {
|
|
cfg, err := r.lidarrCfg.Get(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !cfg.Enabled {
|
|
r.logger.Debug("lidarrrequests: reconciler skipping tick — lidarr disabled")
|
|
return nil
|
|
}
|
|
|
|
q := dbq.New(r.pool)
|
|
rows, err := q.ListApprovedLidarrRequestsForReconcile(ctx, r.batch)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, row := range rows {
|
|
if err := r.reconcileRow(ctx, q, cfg, row); err != nil {
|
|
r.logger.Warn("lidarrrequests: reconciler row failed",
|
|
"request_id", row.ID,
|
|
"kind", row.Kind,
|
|
"err", err,
|
|
)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// reconcileRow drives a request forward by one step. If Lidarr hasn't
|
|
// confirmed the add yet it (re)sends it and returns — the import can't
|
|
// land before Lidarr has the add, so there's nothing to match this tick.
|
|
// Once confirmed, it checks the library for a match and, if found, calls
|
|
// CompleteLidarrRequest. Returns an error only for unexpected failures
|
|
// (the add retry's transient error included, so tickOnce logs WARN and
|
|
// it's retried next tick — "Lidarr just keeps trying"); a no-match is
|
|
// not an error.
|
|
func (r *Reconciler) reconcileRow(ctx context.Context, q *dbq.Queries, cfg lidarrconfig.Config, row dbq.LidarrRequest) error {
|
|
if !row.LidarrAddConfirmedAt.Valid {
|
|
if err := r.ensureLidarrAdd(ctx, q, cfg, row); err != nil {
|
|
return err
|
|
}
|
|
// ensureLidarrAdd either confirmed the add or no-op'd (Lidarr
|
|
// client unavailable). Fall through to import-matching so an
|
|
// item already present in the library still completes this tick.
|
|
}
|
|
switch row.Kind {
|
|
case dbq.LidarrRequestKindArtist:
|
|
return r.reconcileArtist(ctx, q, row)
|
|
case dbq.LidarrRequestKindAlbum:
|
|
return r.reconcileAlbum(ctx, q, row)
|
|
case dbq.LidarrRequestKindTrack:
|
|
return r.reconcileTrack(ctx, q, row)
|
|
default:
|
|
r.logger.Warn("lidarrrequests: reconciler unknown kind", "kind", row.Kind)
|
|
return nil
|
|
}
|
|
}
|
|
|
|
// ensureLidarrAdd (re)sends the Lidarr add for an approved request whose
|
|
// add hasn't been confirmed. On success (including ErrAlreadyExists,
|
|
// handled inside sendLidarrAdd) it stamps lidarr_add_confirmed_at so this
|
|
// stops re-sending and the next tick moves to import-matching. A failure
|
|
// is returned so tickOnce logs WARN and the row is retried next tick;
|
|
// there is intentionally no failed-state or attempt cap — per product
|
|
// decision Lidarr just keeps trying and the operator monitors Lidarr.
|
|
func (r *Reconciler) ensureLidarrAdd(ctx context.Context, q *dbq.Queries, cfg lidarrconfig.Config, row dbq.LidarrRequest) error {
|
|
client := r.clientFn()
|
|
if client == nil {
|
|
// Lidarr disabled mid-flight; tickOnce already gates on
|
|
// cfg.Enabled, so this is just defensive — nothing to do.
|
|
return nil
|
|
}
|
|
err := sendLidarrAdd(ctx, client, cfg, row)
|
|
if errors.Is(err, lidarr.ErrNotFound) {
|
|
// Lidarr's metadata has no album under this id. Requests made before
|
|
// albums carried a release group (#5241) name the album by its
|
|
// MusicBrainz release id, which Lidarr does not index: repoint the
|
|
// request at the release group and try once more.
|
|
if repaired, ok := r.repointAtReleaseGroup(ctx, q, row); ok {
|
|
row = repaired
|
|
err = sendLidarrAdd(ctx, client, cfg, row)
|
|
}
|
|
}
|
|
if err != nil {
|
|
return fmt.Errorf("lidarr add retry: %w", err)
|
|
}
|
|
return q.MarkLidarrRequestAddConfirmed(ctx, row.ID)
|
|
}
|
|
|
|
// repointAtReleaseGroup reads an album or track request's album id as a
|
|
// MusicBrainz release and, when its release group can be named, rewrites the
|
|
// request to it. The library's own albums answer first; MusicBrainz is asked
|
|
// only for a release no album row knows the group of, and its answer is cached
|
|
// onto that album. Reports false when there is nothing to repoint: not an
|
|
// album request, the group is unknown, or the id already is the group.
|
|
func (r *Reconciler) repointAtReleaseGroup(ctx context.Context, q *dbq.Queries, row dbq.LidarrRequest) (dbq.LidarrRequest, bool) {
|
|
if row.Kind == dbq.LidarrRequestKindArtist || row.LidarrAlbumMbid == nil || *row.LidarrAlbumMbid == "" {
|
|
return row, false
|
|
}
|
|
release := *row.LidarrAlbumMbid
|
|
group, err := q.GetReleaseGroupForReleaseMbid(ctx, &release)
|
|
if err != nil {
|
|
if !isNoRows(err) {
|
|
r.logger.Warn("lidarrrequests: release group lookup failed", "request_id", row.ID, "err", err)
|
|
return row, false
|
|
}
|
|
group, err = r.releaseGroup(ctx, release)
|
|
if err != nil {
|
|
if !errors.Is(err, tags.ErrNotFound) {
|
|
r.logger.Warn("lidarrrequests: MusicBrainz release group lookup failed",
|
|
"request_id", row.ID, "release_mbid", release, "err", err)
|
|
}
|
|
return row, false
|
|
}
|
|
if cerr := q.SetReleaseGroupForReleaseMbidIfNull(ctx, dbq.SetReleaseGroupForReleaseMbidIfNullParams{
|
|
ReleaseGroupMbid: group,
|
|
ReleaseMbid: release,
|
|
}); cerr != nil {
|
|
r.logger.Warn("lidarrrequests: cache release group failed", "release_mbid", release, "err", cerr)
|
|
}
|
|
}
|
|
if group == "" || group == release {
|
|
return row, false
|
|
}
|
|
if err := q.SetLidarrRequestAlbumMbid(ctx, dbq.SetLidarrRequestAlbumMbidParams{
|
|
ID: row.ID,
|
|
LidarrAlbumMbid: &group,
|
|
}); err != nil {
|
|
r.logger.Warn("lidarrrequests: repoint request failed", "request_id", row.ID, "err", err)
|
|
return row, false
|
|
}
|
|
r.logger.Info("lidarrrequests: request repointed from release to release group",
|
|
"request_id", row.ID, "release_mbid", release, "release_group_mbid", group)
|
|
row.LidarrAlbumMbid = &group
|
|
return row, true
|
|
}
|
|
|
|
// albumForRequest finds the library album a request's album id names, by
|
|
// release id or release group, and only once it has a track back on disk. A
|
|
// re-acquisition request names an album whose row never went away; matching
|
|
// the row alone would complete the request before Lidarr delivered anything.
|
|
// An exact release match is preferred over another release of the group.
|
|
const albumForRequest = `
|
|
SELECT a.id
|
|
FROM albums a
|
|
WHERE (a.mbid = $1 OR a.release_group_mbid = $1)
|
|
AND EXISTS (SELECT 1 FROM tracks t WHERE t.album_id = a.id AND t.missing_since IS NULL)
|
|
ORDER BY (a.mbid = $1) DESC, a.id
|
|
LIMIT 1`
|
|
|
|
func (r *Reconciler) reconcileArtist(ctx context.Context, q *dbq.Queries, row dbq.LidarrRequest) error {
|
|
var artistID pgtype.UUID
|
|
err := r.pool.QueryRow(ctx,
|
|
"SELECT id FROM artists WHERE mbid = $1",
|
|
row.LidarrArtistMbid,
|
|
).Scan(&artistID)
|
|
if err != nil {
|
|
if isNoRows(err) {
|
|
return nil // not in library yet
|
|
}
|
|
return err
|
|
}
|
|
completed, err := q.CompleteLidarrRequest(ctx, dbq.CompleteLidarrRequestParams{
|
|
ID: row.ID,
|
|
MatchedArtistID: artistID,
|
|
MatchedAlbumID: pgtype.UUID{},
|
|
MatchedTrackID: pgtype.UUID{},
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
r.publishCompleted(completed)
|
|
return nil
|
|
}
|
|
|
|
func (r *Reconciler) reconcileAlbum(ctx context.Context, q *dbq.Queries, row dbq.LidarrRequest) error {
|
|
if row.LidarrAlbumMbid == nil {
|
|
return nil
|
|
}
|
|
var albumID pgtype.UUID
|
|
err := r.pool.QueryRow(ctx, albumForRequest, *row.LidarrAlbumMbid).Scan(&albumID)
|
|
if err != nil {
|
|
if isNoRows(err) {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
completed, err := q.CompleteLidarrRequest(ctx, dbq.CompleteLidarrRequestParams{
|
|
ID: row.ID,
|
|
MatchedAlbumID: albumID,
|
|
MatchedArtistID: pgtype.UUID{},
|
|
MatchedTrackID: pgtype.UUID{},
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
r.publishCompleted(completed)
|
|
return nil
|
|
}
|
|
|
|
func (r *Reconciler) reconcileTrack(ctx context.Context, q *dbq.Queries, row dbq.LidarrRequest) error {
|
|
if row.LidarrAlbumMbid == nil {
|
|
return nil
|
|
}
|
|
// Track-kind requests match via their parent album's MBID, not track.mbid.
|
|
var albumID pgtype.UUID
|
|
err := r.pool.QueryRow(ctx, albumForRequest, *row.LidarrAlbumMbid).Scan(&albumID)
|
|
if err != nil {
|
|
if isNoRows(err) {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
|
|
// Load any track from that album to set matched_track_id.
|
|
var trackID pgtype.UUID
|
|
err = r.pool.QueryRow(ctx,
|
|
"SELECT id FROM tracks WHERE album_id = $1 AND missing_since IS NULL ORDER BY id LIMIT 1",
|
|
albumID,
|
|
).Scan(&trackID)
|
|
if err != nil {
|
|
if isNoRows(err) {
|
|
// Album is in the library but no tracks ingested yet — try again next tick.
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
|
|
completed, err := q.CompleteLidarrRequest(ctx, dbq.CompleteLidarrRequestParams{
|
|
ID: row.ID,
|
|
MatchedAlbumID: albumID,
|
|
MatchedTrackID: trackID,
|
|
MatchedArtistID: pgtype.UUID{},
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
r.publishCompleted(completed)
|
|
return nil
|
|
}
|
|
|
|
func isNoRows(err error) bool {
|
|
return err == pgx.ErrNoRows
|
|
}
|