feat(notifications): requests and flags reach the people who act on them (#5340)
- A new request still pending after any auto-approval notifies the admins (request_pending), but not the requester if they are an admin. A request that dedups into one already in flight is not announced again. - Approving or rejecting a request notifies the requester, and a rejection carries the admin's notes as the reason. An admin deciding their own request gets nothing. - The reconciler notifies the requester when their request arrives (request_completed), linking the matched album or artist. - A request the re-acquisition sweeper files and cannot approve itself notifies the admins. - A quarantine flag notifies every admin except the flagger, naming the track, the flagger and the reason. lidarrrequests.Service.CreateTracked reports whether a request was inserted or deduped; Create wraps it. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
+10
-2
@@ -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
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
+19
-14
@@ -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.
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user