release / govulncheck (push) Successful in 21s
release / web (push) Successful in 1m19s
release / go (push) Successful in 1m39s
release / integration (push) Successful in 5m27s
release / android (push) Successful in 5m47s
release / Build signed APK (releases and dev) (push) Successful in 5m34s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 26s
release / Verify release artifacts (tag releases only) (push) Skipped
Nothing is emailed per event. New music (request_completed) goes out at most once a day, at the summary hour in each user's own timezone, grouped by artist. Everything else is batched: one email a window after the first un-emailed item, holding whatever accumulated. - Migration 0074: notification_email_settings (summary hour, batch window, admin-configurable) and user_notification_email_state (batch start, last sent, failures and retry_after per user and group). Existing rows are stamped emailed so the upgrade sends no backlog. - The Notifier stamps emailed_at at write time when the recipient's email channel is off, so turning email on later doesn't send old items. - Read rows are never selected. A row is stamped only after the mailer accepts, in one transaction with the state, against the read's clock, so a coalesced row updated mid-send stays pending. - A failed send backs off 5m doubling to 6h; SMTP not configured just waits. - Links come from the public address; without one the email has none. - The mailer now RFC 2047-encodes subjects and strips line breaks from them. - Admin → Integrations gains a Notification emails card. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
405 lines
14 KiB
Go
405 lines
14 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/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),
|
|
Artist: row.ArtistName,
|
|
Title: requestedTitle(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
|
|
}
|