diff --git a/cmd/minstrel/main.go b/cmd/minstrel/main.go index a276ee53..f88a5b62 100644 --- a/cmd/minstrel/main.go +++ b/cmd/minstrel/main.go @@ -314,7 +314,13 @@ func run() error { } return lidarr.NewClient(c.BaseURL, c.APIKey) } + // The notifications inbox's one writer (M489), shared by every + // background producer started here. The API builds its own over the + // same pool and bus. + notifier := notifications.New(pool, bus, logger.With("component", "notifications")) + lidarrReconciler := lidarrrequests.NewReconciler(pool, lidarrCfg, lidarrClientFn, logger.With("component", "lidarr"), bus) + lidarrReconciler.SetNotifier(notifier) go lidarrReconciler.Run(ctx) // Missing-file re-acquisition (milestone #290). Turns albums whose files @@ -332,12 +338,14 @@ func run() error { if reacqErr != nil { logger.Warn("reacquisition: using default settings", "err", reacqErr) } - go reacquisition.NewSweeper( + reacqSweeper := reacquisition.NewSweeper( pool, reacqSettings, lidarrrequests.NewService(pool, lidarrCfg, lidarrClientFn, nil), logger.With("component", "reacquisition"), - ).Run(ctx) + ) + reacqSweeper.SetNotifier(notifier) + go reacqSweeper.Run(ctx) // library_changes compactor (#357 follow-up). Daily tick; deletes // rows older than the configured retention so the change-log table diff --git a/internal/api/admin_requests.go b/internal/api/admin_requests.go index 69a61bdb..1d63e487 100644 --- a/internal/api/admin_requests.go +++ b/internal/api/admin_requests.go @@ -12,6 +12,7 @@ import ( "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" "git.fabledsword.com/bvandeusen/minstrel/internal/lidarr" "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrrequests" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" ) // validRequestStatuses is the set of allowed values for the ?status= param. @@ -125,6 +126,7 @@ func (h *handlers) handleApproveRequest(w http.ResponseWriter, r *http.Request) return } + h.notifyRequestDecided(r.Context(), notifications.KindRequestApproved, admin, row, "") h.publishRequestStatusChanged(row) writeJSON(w, http.StatusOK, requestViewFrom(row)) } @@ -165,6 +167,7 @@ func (h *handlers) handleRejectRequest(w http.ResponseWriter, r *http.Request) { return } + h.notifyRequestDecided(r.Context(), notifications.KindRequestRejected, admin, row, body.Notes) h.publishRequestStatusChanged(row) writeJSON(w, http.StatusOK, requestViewFrom(row)) } diff --git a/internal/api/api.go b/internal/api/api.go index 45681b8d..4cffb66f 100644 --- a/internal/api/api.go +++ b/internal/api/api.go @@ -23,6 +23,7 @@ import ( "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrrequests" "git.fabledsword.com/bvandeusen/minstrel/internal/mailer" "git.fabledsword.com/bvandeusen/minstrel/internal/netsettings" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" "git.fabledsword.com/bvandeusen/minstrel/internal/playevents" "git.fabledsword.com/bvandeusen/minstrel/internal/playlists" "git.fabledsword.com/bvandeusen/minstrel/internal/reacquisition" @@ -61,6 +62,7 @@ func Mount(r chi.Router, pool *pgxpool.Pool, logger *slog.Logger, events *playev dataDir: dataDir, mailer: sender, eventbus: bus, + notifier: notifications.New(pool, bus, logger.With("component", "notifications")), playlistScheduler: playlistScheduler, streamSecret: streamSecret, netSettings: netSettings, @@ -333,20 +335,23 @@ type handlers struct { // librarySize memoises the track count that sizes the candidate pool // (#3880). Held here rather than counted per request: the count is a // full table scan, and library size only moves when a scan runs. - librarySize *recommendation.LibrarySize - lidarrCfg *lidarrconfig.Service - lidarrRequests *lidarrrequests.Service - lidarrQuarantine *lidarrquarantine.Service - tracks *tracks.Service - playlists *playlists.Service - coverart *coverart.Enricher - coverSettings *coverart.SettingsService - tagSettings *tags.SettingsService - scanner *library.Scanner - scanCfg library.RunScanConfig - dataDir string - mailer mailer.Sender - eventbus *eventbus.Bus + librarySize *recommendation.LibrarySize + lidarrCfg *lidarrconfig.Service + lidarrRequests *lidarrrequests.Service + lidarrQuarantine *lidarrquarantine.Service + tracks *tracks.Service + playlists *playlists.Service + coverart *coverart.Enricher + coverSettings *coverart.SettingsService + tagSettings *tags.SettingsService + scanner *library.Scanner + scanCfg library.RunScanConfig + dataDir string + mailer mailer.Sender + eventbus *eventbus.Bus + // notifier writes the notifications inbox (M489). Nil-safe: a nil + // notifier records nothing, which is what most handler tests want. + notifier *notifications.Notifier playlistScheduler *playlists.Scheduler // reacqSettings is the DB-backed policy for auto re-acquisition of // missing files (milestone #290) — grace window, backoff, attempt caps. diff --git a/internal/api/auth_test.go b/internal/api/auth_test.go index 4f04a6ef..feb137ae 100644 --- a/internal/api/auth_test.go +++ b/internal/api/auth_test.go @@ -27,6 +27,7 @@ import ( "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrquarantine" "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrrequests" "git.fabledsword.com/bvandeusen/minstrel/internal/mailer" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" "git.fabledsword.com/bvandeusen/minstrel/internal/playevents" "git.fabledsword.com/bvandeusen/minstrel/internal/playlists" "git.fabledsword.com/bvandeusen/minstrel/internal/recsettings" @@ -72,7 +73,7 @@ func testHandlers(t *testing.T) (*handlers, *pgxpool.Pool) { dataDir := t.TempDir() tracksSvc := tracks.NewService(pool, logger, nil, dataDir) playlistsSvc := playlists.NewService(pool, logger, dataDir) - h := &handlers{pool: pool, logger: logger, events: w, recCfg: recCfg, recSettings: recSettings, rng: func() float64 { return 0.5 }, lidarrCfg: lidarrCfg, lidarrRequests: lidarrReqs, lidarrQuarantine: lidarrQuar, tracks: tracksSvc, playlists: playlistsSvc, dataDir: dataDir, scanner: nil, scanCfg: library.RunScanConfig{}, mailer: &mailer.FakeSender{}} + h := &handlers{pool: pool, logger: logger, events: w, recCfg: recCfg, recSettings: recSettings, rng: func() float64 { return 0.5 }, lidarrCfg: lidarrCfg, lidarrRequests: lidarrReqs, lidarrQuarantine: lidarrQuar, tracks: tracksSvc, playlists: playlistsSvc, dataDir: dataDir, scanner: nil, scanCfg: library.RunScanConfig{}, mailer: &mailer.FakeSender{}, notifier: notifications.New(pool, nil, nil)} return h, pool } diff --git a/internal/api/notify_producers.go b/internal/api/notify_producers.go new file mode 100644 index 00000000..b8fe22c2 --- /dev/null +++ b/internal/api/notify_producers.go @@ -0,0 +1,73 @@ +package api + +import ( + "context" + + "github.com/jackc/pgx/v5/pgtype" + + "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" + "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrrequests" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" +) + +// Helpers for the handlers that produce notifications (M489). Every call is +// NotifyLogged: a notification never fails the action that caused it. + +// userLabel is how a user is named to someone else: their display name, or +// their username when they have none. +func userLabel(u dbq.User) string { + if u.DisplayName != nil && *u.DisplayName != "" { + return *u.DisplayName + } + return u.Username +} + +func requestPayload(row dbq.LidarrRequest) notifications.Payload { + return notifications.Payload{ + RequestID: uuidToString(row.ID), + RequestKind: string(row.Kind), + Name: lidarrrequests.DisplayName(row), + } +} + +// notifyRequestDecided tells the requester an admin approved or declined +// their request. An admin deciding their own request needs no notice. +func (h *handlers) notifyRequestDecided(ctx context.Context, kind notifications.Kind, admin dbq.User, row dbq.LidarrRequest, reason string) { + if row.UserID == admin.ID { + return + } + p := requestPayload(row) + p.Reason = reason + h.notifier.NotifyLogged(ctx, kind, notifications.ToUser(row.UserID), p.Map()) +} + +// quarantineReasonLabel is a flag reason as the flag popover words it, in +// lower case to sit inside a sentence ("flagged X: bad rip"). +func quarantineReasonLabel(reason string) string { + switch reason { + case "bad_rip": + return "bad rip" + case "wrong_file": + return "wrong file" + case "wrong_tags": + return "wrong tags" + case "duplicate": + return "duplicate" + default: + return "" + } +} + +// trackLabel names a track "Artist – Title" for a notification, falling back +// to the title alone, or to nothing, rather than failing the caller. +func (h *handlers) trackLabel(ctx context.Context, trackID pgtype.UUID) string { + q := dbq.New(h.pool) + t, err := q.GetTrackByID(ctx, trackID) + if err != nil { + return "" + } + if a, aerr := q.GetArtistByID(ctx, t.ArtistID); aerr == nil && a.Name != "" { + return a.Name + " – " + t.Title + } + return t.Title +} diff --git a/internal/api/notify_producers_test.go b/internal/api/notify_producers_test.go new file mode 100644 index 00000000..d90db858 --- /dev/null +++ b/internal/api/notify_producers_test.go @@ -0,0 +1,212 @@ +package api + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/jackc/pgx/v5/pgxpool" + + "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" +) + +// inboxItem is one thing the notifier left for one user: kind and decoded payload, +// newest first. +type inboxItem struct { + Kind string + Payload notifications.Payload +} + +func inboxOf(t *testing.T, pool *pgxpool.Pool, user dbq.User) []inboxItem { + t.Helper() + rows, err := dbq.New(pool).ListNotifications(context.Background(), dbq.ListNotificationsParams{ + UserID: user.ID, PageLimit: 100, + }) + if err != nil { + t.Fatalf("list notifications: %v", err) + } + out := make([]inboxItem, 0, len(rows)) + for _, r := range rows { + var p notifications.Payload + if err := json.Unmarshal(r.Payload, &p); err != nil { + t.Fatalf("payload %s: %v", r.Payload, err) + } + out = append(out, inboxItem{Kind: r.Kind, Payload: p}) + } + return out +} + +func wantInbox(t *testing.T, who string, got []inboxItem, kinds ...notifications.Kind) { + t.Helper() + if len(got) != len(kinds) { + t.Fatalf("%s inbox = %+v, want kinds %v", who, got, kinds) + } + for i, k := range kinds { + if got[i].Kind != string(k) { + t.Errorf("%s inbox[%d] = %q, want %q", who, i, got[i].Kind, k) + } + } +} + +// newApprovingLidarrStub answers every Lidarr call an approval makes. +func newApprovingLidarrStub(t *testing.T) *httptest.Server { + t.Helper() + stub := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + switch { + case strings.Contains(r.URL.Path, "/metadataprofile"): + _, _ = w.Write([]byte(`[{"id":1,"name":"Standard"}]`)) + case strings.Contains(r.URL.Path, "/qualityprofile"): + _, _ = w.Write([]byte(`[{"id":1,"name":"Lossless"}]`)) + default: + _, _ = w.Write([]byte(`{"id":1}`)) + } + })) + t.Cleanup(stub.Close) + return stub +} + +func TestNotify_NewRequestReachesOtherAdminsOnce(t *testing.T) { + h, pool := testHandlers(t) + resetLidarrState(t, h) + + alice := seedUser(t, pool, "np-alice", "pw", false) + admin := seedUser(t, pool, "np-admin", "pw", true) + selfAdmin := seedUser(t, pool, "np-self", "pw", true) + + rv := createArtistRequest(t, h, alice, "np-mbid-1", "Pending Band") + // The same request again dedups into the first and is not announced twice. + createArtistRequest(t, h, alice, "np-mbid-1", "Pending Band") + + got := inboxOf(t, pool, admin) + wantInbox(t, "admin", got, notifications.KindRequestPending) + if got[0].Payload.RequestID != uuidToString(rv.ID) || got[0].Payload.Name != "Pending Band" || got[0].Payload.Actor != alice.Username { + t.Errorf("pending payload = %+v", got[0].Payload) + } + wantInbox(t, "requester", inboxOf(t, pool, alice)) + + // An admin's own request goes to the other admins, never back to them. + createArtistRequest(t, h, selfAdmin, "np-mbid-2", "Own Band") + wantInbox(t, "self-admin", inboxOf(t, pool, selfAdmin), notifications.KindRequestPending) + if n := len(inboxOf(t, pool, admin)); n != 2 { + t.Errorf("admin inbox = %d, want 2 (alice's and the self-admin's)", n) + } +} + +func TestNotify_AutoApprovedRequestWaitsOnNobody(t *testing.T) { + h, _ := testHandlersWithClientFn(t) + resetLidarrState(t, h) + saveLidarrConfig(t, h, newApprovingLidarrStub(t).URL, true) + + alice := seedUser(t, h.pool, "np-auto", "pw", false) + admin := seedUser(t, h.pool, "np-auto-admin", "pw", true) + if _, err := h.pool.Exec(context.Background(), + "UPDATE users SET auto_approve_requests = true WHERE id = $1", alice.ID); err != nil { + t.Fatalf("set auto_approve: %v", err) + } + alice, err := dbq.New(h.pool).GetUserByID(context.Background(), alice.ID) + if err != nil { + t.Fatalf("reload user: %v", err) + } + + rv := createArtistRequest(t, h, alice, "np-auto-mbid", "Auto Band") + if rv.Status != "approved" { + t.Fatalf("status = %q, want approved — the stub should accept the add", rv.Status) + } + wantInbox(t, "admin", inboxOf(t, h.pool, admin)) + wantInbox(t, "requester", inboxOf(t, h.pool, alice)) +} + +func TestNotify_DecisionsReachTheRequester(t *testing.T) { + h, _ := testHandlersWithClientFn(t) + resetLidarrState(t, h) + saveLidarrConfig(t, h, newApprovingLidarrStub(t).URL, true) + + alice := seedUser(t, h.pool, "np-dec-alice", "pw", false) + admin := seedUser(t, h.pool, "np-dec-admin", "pw", true) + + approved := seedPendingArtistRequest(t, h, alice, "np-dec-1", "Yes Band") + rejected := seedPendingArtistRequest(t, h, alice, "np-dec-2", "No Band") + own := seedPendingArtistRequest(t, h, admin, "np-dec-3", "Own Band") + + for _, c := range []struct { + id string + verb string + body []byte + }{ + {uuidToString(approved.ID), "approve", nil}, + {uuidToString(rejected.ID), "reject", []byte(`{"notes":"Only a live bootleg exists"}`)}, + {uuidToString(own.ID), "approve", nil}, + } { + w := doAdminRequestReq(t, h, http.MethodPost, "/api/admin/requests/"+c.id+"/"+c.verb, c.body, admin) + if w.Code != http.StatusOK { + t.Fatalf("%s %s: status = %d; body = %s", c.verb, c.id, w.Code, w.Body.String()) + } + } + + got := inboxOf(t, h.pool, alice) + wantInbox(t, "requester", got, notifications.KindRequestRejected, notifications.KindRequestApproved) + if got[0].Payload.Name != "No Band" || got[0].Payload.Reason != "Only a live bootleg exists" { + t.Errorf("rejected payload = %+v", got[0].Payload) + } + if got[1].Payload.Name != "Yes Band" || got[1].Payload.Reason != "" { + t.Errorf("approved payload = %+v", got[1].Payload) + } + // The admin decided their own request: no notice. alice's two pending + // notices are all the admin holds. + for _, it := range inboxOf(t, h.pool, admin) { + if it.Kind != string(notifications.KindRequestPending) { + t.Errorf("admin inbox holds %q, want only request_pending", it.Kind) + } + } +} + +func TestNotify_FlagReachesAdminsButNotTheFlagger(t *testing.T) { + h, pool := testHandlers(t) + truncateLibrary(t, pool) + + alice := seedUser(t, pool, "np-flag-alice", "pw", false) + admin := seedUser(t, pool, "np-flag-admin", "pw", true) + other := seedUser(t, pool, "np-flag-other", "pw", true) + track := seedQuarantineTrack(t, h, "np") + + w := doFlag(h, alice, fmt.Sprintf(`{"track_id":%q,"reason":"bad_rip"}`, uuidToString(track.ID))) + if w.Code != http.StatusCreated { + t.Fatalf("flag: status = %d; body = %s", w.Code, w.Body.String()) + } + got := inboxOf(t, pool, admin) + wantInbox(t, "admin", got, notifications.KindQuarantineFlagged) + if p := got[0].Payload; p.Name != "Q Artist np – Q Track np" || p.Actor != alice.Username || p.Reason != "bad rip" { + t.Errorf("flag payload = %+v", p) + } + wantInbox(t, "flagger", inboxOf(t, pool, alice)) + + // An admin flagging a track tells the other admins only. + track2 := seedQuarantineTrack(t, h, "np2") + w = doFlag(h, admin, fmt.Sprintf(`{"track_id":%q,"reason":"other"}`, uuidToString(track2.ID))) + if w.Code != http.StatusCreated { + t.Fatalf("admin flag: status = %d; body = %s", w.Code, w.Body.String()) + } + if n := len(inboxOf(t, pool, admin)); n != 1 { + t.Errorf("flagging admin inbox = %d, want 1 (alice's flag only)", n) + } + if n := len(inboxOf(t, pool, other)); n != 2 { + t.Errorf("other admin inbox = %d, want 2", n) + } +} + +func TestQuarantineReasonLabel(t *testing.T) { + for in, want := range map[string]string{ + "bad_rip": "bad rip", "wrong_file": "wrong file", "wrong_tags": "wrong tags", + "duplicate": "duplicate", "other": "", "": "", + } { + if got := quarantineReasonLabel(in); got != want { + t.Errorf("quarantineReasonLabel(%q) = %q, want %q", in, got, want) + } + } +} diff --git a/internal/api/quarantine.go b/internal/api/quarantine.go index 1beec8b2..c5df1b22 100644 --- a/internal/api/quarantine.go +++ b/internal/api/quarantine.go @@ -9,6 +9,7 @@ import ( "git.fabledsword.com/bvandeusen/minstrel/internal/apierror" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrquarantine" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" ) // quarantineView is the JSON shape returned by the user-facing endpoints @@ -70,6 +71,10 @@ func (h *handlers) handleFlag(w http.ResponseWriter, r *http.Request) { // Broadcast: the flagging user's other clients invalidate their // Hidden tab; admins' clients invalidate their quarantine queue. h.publishQuarantineEvent("quarantine.flagged", user.ID, trackID, true) + // Admins review flags (M489). An admin flagging a track needs no notice + // of their own flag. + h.notifier.NotifyLogged(r.Context(), notifications.KindQuarantineFlagged, notifications.ToAdmins(user.ID), + notifications.Payload{Name: h.trackLabel(r.Context(), trackID), Actor: userLabel(user), Reason: quarantineReasonLabel(body.Reason)}.Map()) writeJSON(w, http.StatusCreated, quarantineViewFrom(row)) } diff --git a/internal/api/requests.go b/internal/api/requests.go index fcd84b95..0c1edd47 100644 --- a/internal/api/requests.go +++ b/internal/api/requests.go @@ -10,6 +10,7 @@ import ( "git.fabledsword.com/bvandeusen/minstrel/internal/apierror" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrrequests" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" ) // requestView is the JSON shape returned by all /api/requests handlers. @@ -130,7 +131,7 @@ func (h *handlers) handleCreateRequest(w http.ResponseWriter, r *http.Request) { return } - row, err := h.lidarrRequests.Create(r.Context(), user.ID, lidarrrequests.CreateParams{ + row, created, err := h.lidarrRequests.CreateTracked(r.Context(), user.ID, lidarrrequests.CreateParams{ Kind: body.Kind, LidarrArtistMBID: body.LidarrArtistMBID, LidarrAlbumMBID: body.LidarrAlbumMBID, @@ -175,6 +176,16 @@ func (h *handlers) handleCreateRequest(w http.ResponseWriter, r *http.Request) { } } + // A new request still pending after any auto-approval waits on an admin + // (M489). One that deduped into a request already in flight was + // announced when it was first made, and the requester, if they are an + // admin themselves, needs no notice of their own request. + if created && row.Status == dbq.LidarrRequestStatusPending { + p := requestPayload(row) + p.Actor = userLabel(user) + h.notifier.NotifyLogged(r.Context(), notifications.KindRequestPending, notifications.ToAdmins(user.ID), p.Map()) + } + h.publishRequestStatusChanged(row) writeJSON(w, http.StatusCreated, requestViewFrom(row)) } diff --git a/internal/lidarrrequests/reconciler.go b/internal/lidarrrequests/reconciler.go index 122bdb4d..1f729d6e 100644 --- a/internal/lidarrrequests/reconciler.go +++ b/internal/lidarrrequests/reconciler.go @@ -15,6 +15,7 @@ import ( "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" ) @@ -30,6 +31,7 @@ type Reconciler struct { 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 @@ -61,9 +63,26 @@ func NewReconciler(pool *pgxpool.Pool, cfg *lidarrconfig.Service, clientFn func( } } -// 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) { +// 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 } @@ -308,7 +327,7 @@ func (r *Reconciler) reconcileArtist(ctx context.Context, q *dbq.Queries, row db if err != nil { return err } - r.publishCompleted(completed) + r.publishCompleted(ctx, completed) return nil } @@ -333,7 +352,7 @@ func (r *Reconciler) reconcileAlbum(ctx context.Context, q *dbq.Queries, row dbq if err != nil { return err } - r.publishCompleted(completed) + r.publishCompleted(ctx, completed) return nil } @@ -374,7 +393,7 @@ func (r *Reconciler) reconcileTrack(ctx context.Context, q *dbq.Queries, row dbq if err != nil { return err } - r.publishCompleted(completed) + r.publishCompleted(ctx, completed) return nil } diff --git a/internal/lidarrrequests/reconciler_notify_test.go b/internal/lidarrrequests/reconciler_notify_test.go new file mode 100644 index 00000000..53e71751 --- /dev/null +++ b/internal/lidarrrequests/reconciler_notify_test.go @@ -0,0 +1,52 @@ +package lidarrrequests + +import ( + "context" + "encoding/json" + "testing" + + "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" + "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrconfig" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" +) + +// M489: an album request that arrives lands in the requester's inbox once, +// linked to the album it matched. +func TestReconciler_CompletionNotifiesRequester(t *testing.T) { + pool := newPool(t) + q := dbq.New(pool) + ctx := context.Background() + + enableLidarrForPool(t, pool) + user := seedUser(t, pool) + artist := seedArtist(t, q, "Notify Artist", "notify-artist-mbid") + album := seedAlbum(t, q, artist.ID, "Notify Album", "notify-album-mbid") + _ = seedTrack(t, q, album.ID, artist.ID, "Notify Track", "/music/notify/01.flac") + seedApprovedRequestDirect(t, q, user, CreateParams{ + Kind: "album", LidarrArtistMBID: "notify-artist-mbid", LidarrAlbumMBID: "notify-album-mbid", + ArtistName: "Notify Artist", AlbumTitle: "Notify Album", + }) + + rec := NewReconciler(pool, lidarrconfig.New(pool), nil, newTestLogger(), nil) + rec.SetNotifier(notifications.New(pool, nil, nil)) + for i := 0; i < 2; i++ { // the second tick finds nothing left to complete + if err := rec.tickOnce(ctx); err != nil { + t.Fatalf("tickOnce: %v", err) + } + } + + rows, err := q.ListNotifications(ctx, dbq.ListNotificationsParams{UserID: user, PageLimit: 10}) + if err != nil { + t.Fatalf("list notifications: %v", err) + } + if len(rows) != 1 || rows[0].Kind != string(notifications.KindRequestCompleted) { + t.Fatalf("inbox = %+v, want one request_completed", rows) + } + var p notifications.Payload + if err := json.Unmarshal(rows[0].Payload, &p); err != nil { + t.Fatal(err) + } + if p.Name != "Notify Artist – Notify Album" || p.AlbumID != formatUUIDForBus(album.ID) || p.ArtistID != "" { + t.Errorf("payload = %+v", p) + } +} diff --git a/internal/lidarrrequests/service.go b/internal/lidarrrequests/service.go index c5721732..bf37ff28 100644 --- a/internal/lidarrrequests/service.go +++ b/internal/lidarrrequests/service.go @@ -80,8 +80,16 @@ func NewService(pool *pgxpool.Pool, cfg *lidarrconfig.Service, clientFn func() * // Create validates the kind→required-fields invariant and inserts a // pending row. func (s *Service) Create(ctx context.Context, userID pgtype.UUID, p CreateParams) (dbq.LidarrRequest, error) { + row, _, err := s.CreateTracked(ctx, userID, p) + return row, err +} + +// CreateTracked is Create, also reporting whether a row was inserted. A +// request that dedupes into one already in flight is not new, so callers +// that announce new requests (M489) say nothing for it. +func (s *Service) CreateTracked(ctx context.Context, userID pgtype.UUID, p CreateParams) (dbq.LidarrRequest, bool, error) { if err := validateKindFields(p); err != nil { - return dbq.LidarrRequest{}, err + return dbq.LidarrRequest{}, false, err } q := dbq.New(s.pool) // Dedup: if a non-terminal request for this MBID already exists, @@ -97,9 +105,9 @@ func (s *Service) Create(ctx context.Context, userID pgtype.UUID, p CreateParams dedupMBID = p.LidarrTrackMBID } if existing, derr := q.GetNonTerminalRequestForMBID(ctx, dedupMBID); derr == nil { - return existing, nil + return existing, false, nil } else if !errors.Is(derr, pgx.ErrNoRows) { - return dbq.LidarrRequest{}, fmt.Errorf("create: dedup check: %w", derr) + return dbq.LidarrRequest{}, false, fmt.Errorf("create: dedup check: %w", derr) } row, err := q.CreateLidarrRequest(ctx, dbq.CreateLidarrRequestParams{ UserID: userID, @@ -112,9 +120,25 @@ func (s *Service) Create(ctx context.Context, userID pgtype.UUID, p CreateParams TrackTitle: strPtr(p.TrackTitle), }) if err != nil { - return dbq.LidarrRequest{}, fmt.Errorf("create: %w", err) + return dbq.LidarrRequest{}, false, fmt.Errorf("create: %w", err) } - return row, nil + return row, true, nil +} + +// DisplayName is a request as a person would say it: "Moe Shop" for an +// artist, "Moe Shop – WWW" for an album or a track. +func DisplayName(row dbq.LidarrRequest) string { + switch row.Kind { + case dbq.LidarrRequestKindTrack: + if row.TrackTitle != nil && *row.TrackTitle != "" { + return row.ArtistName + " – " + *row.TrackTitle + } + case dbq.LidarrRequestKindAlbum: + if row.AlbumTitle != nil && *row.AlbumTitle != "" { + return row.ArtistName + " – " + *row.AlbumTitle + } + } + return row.ArtistName } func validateKindFields(p CreateParams) error { diff --git a/internal/reacquisition/sweeper.go b/internal/reacquisition/sweeper.go index 16875b40..fba2a6e3 100644 --- a/internal/reacquisition/sweeper.go +++ b/internal/reacquisition/sweeper.go @@ -7,12 +7,14 @@ import ( "log/slog" "time" + "github.com/google/uuid" "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/lidarrrequests" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" "git.fabledsword.com/bvandeusen/minstrel/internal/tags" ) @@ -20,7 +22,7 @@ import ( // narrowed to an interface so the pass can be tested without a Lidarr client // or an approval path that talks to one. type requestCreator interface { - Create(ctx context.Context, userID pgtype.UUID, p lidarrrequests.CreateParams) (dbq.LidarrRequest, error) + CreateTracked(ctx context.Context, userID pgtype.UUID, p lidarrrequests.CreateParams) (dbq.LidarrRequest, bool, error) Approve(ctx context.Context, requestID, adminID pgtype.UUID, ov lidarrrequests.ApproveOverrides) (dbq.LidarrRequest, error) } @@ -37,6 +39,9 @@ type Sweeper struct { settings *SettingsService requests requestCreator logger *slog.Logger + // notifier tells admins about a request the sweeper filed that still + // needs their approval (M489). Nil: nobody is told. + notifier *notifications.Notifier tick time.Duration // releaseGroup names the MusicBrainz release group of a release id, for an // album whose tags never carried one (#5241). Lidarr knows albums only by @@ -92,6 +97,10 @@ type PassResult struct { Unresolved int } +// SetNotifier makes requests the sweeper files, and cannot approve itself, +// reach the admins' notifications inbox (M489). +func (s *Sweeper) SetNotifier(n *notifications.Notifier) { s.notifier = n } + // SweepOnce runs one pass. Exported so the admin surface can offer a "run // now" without waiting out the tick, and so tests drive it directly. func (s *Sweeper) SweepOnce(ctx context.Context) error { @@ -182,7 +191,7 @@ func (s *Sweeper) attempt( return fmt.Errorf("release group: %w", err) } - req, err := s.requests.Create(ctx, adminID, lidarrrequests.CreateParams{ + req, created, err := s.requests.CreateTracked(ctx, adminID, lidarrrequests.CreateParams{ Kind: "album", LidarrArtistMBID: *album.ArtistMbid, ArtistName: album.ArtistName, @@ -206,11 +215,15 @@ func (s *Sweeper) attempt( return fmt.Errorf("record attempt: %w", err) } + // A request the sweeper filed and could not approve waits on an admin. + // One it deduped into was announced when it was first filed. + awaitingAdmin := created if cfg.AutoApprove { _, aerr := s.requests.Approve(ctx, req.ID, adminID, lidarrrequests.ApproveOverrides{}) switch { case aerr == nil: res.Approved++ + awaitingAdmin = false case errors.Is(aerr, lidarrrequests.ErrLidarrDisabled): // Leave it pending rather than treating it as a failure. The // request is still the right record of intent, and it becomes @@ -225,6 +238,16 @@ func (s *Sweeper) attempt( } } + if awaitingAdmin { + s.notifier.NotifyLogged(ctx, notifications.KindRequestPending, notifications.ToAdmins(pgtype.UUID{}), + notifications.Payload{ + RequestID: uuid.UUID(req.ID.Bytes).String(), + RequestKind: string(req.Kind), + Name: lidarrrequests.DisplayName(req), + Actor: "Re-acquisition", + }.Map()) + } + if row.Attempts >= cfg.MaxAttempts { if err := q.MarkReacquisitionGaveUp(ctx, album.AlbumID); err != nil { return fmt.Errorf("mark gave up: %w", err) diff --git a/internal/reacquisition/sweeper_test.go b/internal/reacquisition/sweeper_test.go index d8833a69..0bb7c732 100644 --- a/internal/reacquisition/sweeper_test.go +++ b/internal/reacquisition/sweeper_test.go @@ -15,6 +15,7 @@ import ( "git.fabledsword.com/bvandeusen/minstrel/internal/dbtest" "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrconfig" "git.fabledsword.com/bvandeusen/minstrel/internal/lidarrrequests" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" "git.fabledsword.com/bvandeusen/minstrel/internal/tags" ) @@ -84,9 +85,10 @@ func TestSweep_RequestsAlbumsByReleaseGroup_Integration(t *testing.T) { logger := slog.New(slog.NewTextHandler(io.Discard, nil)) q := dbq.New(pool) - if _, err := q.CreateUser(ctx, dbq.CreateUserParams{ + admin, err := q.CreateUser(ctx, dbq.CreateUserParams{ Username: dbtest.TestUserPrefix + "sweepadmin", PasswordHash: "x", ApiTokenHash: "x", IsAdmin: true, - }); err != nil { + }) + if err != nil { t.Fatal(err) } settings, err := NewSettingsService(ctx, pool, logger) @@ -105,6 +107,7 @@ func TestSweep_RequestsAlbumsByReleaseGroup_Integration(t *testing.T) { requests := lidarrrequests.NewService(pool, lidarrconfig.New(pool), nil, nil) s := NewSweeper(pool, settings, requests, logger) + s.SetNotifier(notifications.New(pool, nil, nil)) asked := map[string]int{} s.releaseGroup = func(_ context.Context, release string) (string, error) { asked[release]++ @@ -153,6 +156,20 @@ SELECT lr.lidarr_album_mbid if attempts != 0 { t.Errorf("unresolvable album spent an attempt (%d state rows)", attempts) } + // M489: the two requests filed without auto-approval wait on an admin, + // so each reaches the admins' inbox; the unresolvable album filed none. + inbox, err := q.ListNotifications(ctx, dbq.ListNotificationsParams{UserID: admin.ID, PageLimit: 10}) + if err != nil { + t.Fatal(err) + } + if len(inbox) != 2 { + t.Errorf("admin inbox holds %d notifications, want 2", len(inbox)) + } + for _, n := range inbox { + if n.Kind != string(notifications.KindRequestPending) { + t.Errorf("admin inbox holds %q, want request_pending", n.Kind) + } + } var stray int if err := pool.QueryRow(ctx, "SELECT count(*) FROM lidarr_requests WHERE lidarr_album_mbid LIKE 'release-%'").Scan(&stray); err != nil { t.Fatal(err)