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/notifications" "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 notifier *notifications.Notifier // nil: no inbox notifications 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, } } // SetNotifier makes completions land in the requester's notifications inbox // (M489). Without one, completions are only broadcast on the bus. func (r *Reconciler) SetNotifier(n *notifications.Notifier) { r.notifier = n } // publishCompleted tells the requester their request arrived: a // request.status_changed event for open screens, and a notification that // outlives the connection. func (r *Reconciler) publishCompleted(ctx context.Context, row dbq.LidarrRequest) { p := notifications.Payload{ RequestID: formatUUIDForBus(row.ID), RequestKind: string(row.Kind), Name: DisplayName(row), } if row.MatchedAlbumID.Valid { p.AlbumID = formatUUIDForBus(row.MatchedAlbumID) } else if row.MatchedArtistID.Valid { p.ArtistID = formatUUIDForBus(row.MatchedArtistID) } r.notifier.NotifyLogged(ctx, notifications.KindRequestCompleted, notifications.ToUser(row.UserID), p.Map()) 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 ($1) names, by // release id or release group, once the request ($2 = requested_at) has been // answered: the album has a track on disk AND either // // - a track arrived after the request (a new album, or Lidarr fetching // another release of the group into its own row), or // - no track that was already missing when the request was made is still // missing (the lost files came back in place or were adopted at a new path). // // "Has a track on disk" alone is not enough (#5263): re-acquisition targets // albums with ANY track missing, so the tracks that never left satisfied it // the moment Lidarr accepted the add. added_at is the arrival clock because // nothing rewrites it; updated_at moves on every tag re-read. // 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) AND ( EXISTS (SELECT 1 FROM tracks t WHERE t.album_id = a.id AND t.missing_since IS NULL AND t.added_at > $2) OR NOT EXISTS (SELECT 1 FROM tracks t WHERE t.album_id = a.id AND t.missing_since IS NOT NULL AND t.missing_since < $2) ) 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(ctx, 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, row.RequestedAt).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(ctx, 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, row.RequestedAt).Scan(&albumID) if err != nil { if isNoRows(err) { return nil } return err } // Load a track from that album to set matched_track_id, newest arrival first. var trackID pgtype.UUID err = r.pool.QueryRow(ctx, "SELECT id FROM tracks WHERE album_id = $1 AND missing_since IS NULL ORDER BY added_at DESC, 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(ctx, completed) return nil } func isNoRows(err error) bool { return err == pgx.ErrNoRows }