feat(notifications): library health reaches admins, coalesced (#5341)
release / govulncheck (push) Successful in 18s
release / web (push) Successful in 1m22s
release / go (push) Successful in 1m43s
release / integration (push) Successful in 4m56s
release / android (push) Successful in 5m16s
release / Build signed APK (releases and dev) (push) Successful in 5m27s
release / Attach APK to the Release (tag releases only) (push) Skipped
release / Build + push container image (push) Successful in 15s
release / Verify release artifacts (tag releases only) (push) Skipped

- A failed scan run sends scan_failed. Each failure adds to the count and
  the notice shows the latest error. A scan cut short by shutdown says
  nothing.
- Marking tracks missing sends tracks_missing with a running count.
- A duplicate sweep that proposes a group it had not proposed before
  sends duplicates_found, counting everything awaiting review. A sweep
  that only re-finds known groups stays quiet, so a read notice isn't
  repeated every sweep (CountDuplicateGroupsDetectedSince).
- A playback-error report sends playback_errors, counting the unresolved
  errors (CountUnresolvedPlaybackErrors).

The library package gets its notifier as a package-level SetNotifier
beside SetEventBus, for the same reason the bus is package-level.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-10-08 07:20:33 -04:00
co-authored by Claude Opus 5.5
parent 069baeb14d
commit 1e9408c835
13 changed files with 338 additions and 0 deletions
+1
View File
@@ -321,6 +321,7 @@ func run() error {
lidarrReconciler := lidarrrequests.NewReconciler(pool, lidarrCfg, lidarrClientFn, logger.With("component", "lidarr"), bus) lidarrReconciler := lidarrrequests.NewReconciler(pool, lidarrCfg, lidarrClientFn, logger.With("component", "lidarr"), bus)
lidarrReconciler.SetNotifier(notifier) lidarrReconciler.SetNotifier(notifier)
library.SetNotifier(notifier)
go lidarrReconciler.Run(ctx) go lidarrReconciler.Run(ctx)
// Missing-file re-acquisition (milestone #290). Turns albums whose files // Missing-file re-acquisition (milestone #290). Turns albums whose files
+13
View File
@@ -71,3 +71,16 @@ func (h *handlers) trackLabel(ctx context.Context, trackID pgtype.UUID) string {
} }
return t.Title return t.Title
} }
// notifyPlaybackErrors tells admins how many playback errors await review.
// The notice coalesces, so while it is unread a new report only updates the
// count.
func (h *handlers) notifyPlaybackErrors(ctx context.Context) {
n, err := dbq.New(h.pool).CountUnresolvedPlaybackErrors(ctx)
if err != nil {
h.logger.Warn("api: playback_error: count unresolved", "err", err)
return
}
h.notifier.NotifyLogged(ctx, notifications.KindPlaybackErrors, notifications.ToAdmins(pgtype.UUID{}),
notifications.Payload{Count: n}.Map())
}
+26
View File
@@ -210,3 +210,29 @@ func TestQuarantineReasonLabel(t *testing.T) {
} }
} }
} }
func TestNotify_PlaybackErrorsCoalesceIntoOneCount(t *testing.T) {
h, pool := testHandlers(t)
truncateLibrary(t, pool)
alice := seedUser(t, pool, "np-play-alice", "pw", false)
admin := seedUser(t, pool, "np-play-admin", "pw", true)
track := seedQuarantineTrack(t, h, "play")
for i := 0; i < 2; i++ {
body := fmt.Sprintf(`{"track_id":%q,"kind":"load_failed","client_id":"web-%d"}`, uuidToString(track.ID), i)
req := withUser(httptest.NewRequest(http.MethodPost, "/api/playback-errors", strings.NewReader(body)), alice)
w := httptest.NewRecorder()
h.handleReportPlaybackError(w, req)
if w.Code != http.StatusCreated {
t.Fatalf("report %d: status = %d; body = %s", i, w.Code, w.Body.String())
}
}
got := inboxOf(t, pool, admin)
wantInbox(t, "admin", got, notifications.KindPlaybackErrors)
if got[0].Payload.Count != 2 {
t.Errorf("count = %d, want 2 unresolved", got[0].Payload.Count)
}
wantInbox(t, "reporter", inboxOf(t, pool, alice))
}
+1
View File
@@ -130,6 +130,7 @@ func (h *handlers) handleReportPlaybackError(w http.ResponseWriter, r *http.Requ
writeErr(w, apierror.Internal(err)) writeErr(w, apierror.Internal(err))
return return
} }
h.notifyPlaybackErrors(r.Context())
writeJSON(w, http.StatusCreated, map[string]string{"id": uuidToString(row.ID)}) writeJSON(w, http.StatusCreated, map[string]string{"id": uuidToString(row.ID)})
} }
+18
View File
@@ -27,6 +27,24 @@ func (q *Queries) AddDuplicateGroupMember(ctx context.Context, arg AddDuplicateG
return err return err
} }
const countDuplicateGroupsDetectedSince = `-- name: CountDuplicateGroupsDetectedSince :one
SELECT count(*)::bigint
FROM duplicate_groups
WHERE status = 'pending'
AND detected_at >= $1
`
// Pending proposals first made at or after `since`: what a sweep that started
// then found for the first time. A refreshed proposal keeps its detected_at,
// so a sweep that only re-finds known groups counts none (M489: admins are
// told about new duplicates, not reminded of the same ones every sweep).
func (q *Queries) CountDuplicateGroupsDetectedSince(ctx context.Context, since pgtype.Timestamptz) (int64, error) {
row := q.db.QueryRow(ctx, countDuplicateGroupsDetectedSince, since)
var column_1 int64
err := row.Scan(&column_1)
return column_1, err
}
const countPendingDuplicateGroups = `-- name: CountPendingDuplicateGroups :one const countPendingDuplicateGroups = `-- name: CountPendingDuplicateGroups :one
SELECT count(*)::bigint SELECT count(*)::bigint
FROM duplicate_groups g FROM duplicate_groups g
+11
View File
@@ -11,6 +11,17 @@ import (
"github.com/jackc/pgx/v5/pgtype" "github.com/jackc/pgx/v5/pgtype"
) )
const countUnresolvedPlaybackErrors = `-- name: CountUnresolvedPlaybackErrors :one
SELECT count(*)::bigint FROM playback_errors WHERE resolved_at IS NULL
`
func (q *Queries) CountUnresolvedPlaybackErrors(ctx context.Context) (int64, error) {
row := q.db.QueryRow(ctx, countUnresolvedPlaybackErrors)
var column_1 int64
err := row.Scan(&column_1)
return column_1, err
}
const insertPlaybackError = `-- name: InsertPlaybackError :one const insertPlaybackError = `-- name: InsertPlaybackError :one
INSERT INTO playback_errors (track_id, user_id, client_id, kind, detail) INSERT INTO playback_errors (track_id, user_id, client_id, kind, detail)
VALUES ($1, $2, $3, $4, $5) VALUES ($1, $2, $3, $4, $5)
+10
View File
@@ -154,3 +154,13 @@ SELECT p.id AS group_id,
UPDATE duplicate_groups UPDATE duplicate_groups
SET status = 'dismissed', resolved_at = now() SET status = 'dismissed', resolved_at = now()
WHERE id = sqlc.arg(id) AND status = 'pending'; WHERE id = sqlc.arg(id) AND status = 'pending';
-- name: CountDuplicateGroupsDetectedSince :one
-- Pending proposals first made at or after `since`: what a sweep that started
-- then found for the first time. A refreshed proposal keeps its detected_at,
-- so a sweep that only re-finds known groups counts none (M489: admins are
-- told about new duplicates, not reminded of the same ones every sweep).
SELECT count(*)::bigint
FROM duplicate_groups
WHERE status = 'pending'
AND detected_at >= sqlc.arg(since);
+3
View File
@@ -61,3 +61,6 @@ UPDATE playback_errors
resolved_by = $2, resolved_by = $2,
resolution = $3 resolution = $3
WHERE track_id = $1 AND resolved_at IS NULL; WHERE track_id = $1 AND resolved_at IS NULL;
-- name: CountUnresolvedPlaybackErrors :one
SELECT count(*)::bigint FROM playback_errors WHERE resolved_at IS NULL;
+26
View File
@@ -14,6 +14,7 @@ import (
"github.com/jackc/pgx/v5/pgxpool" "github.com/jackc/pgx/v5/pgxpool"
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
"git.fabledsword.com/bvandeusen/minstrel/internal/notifications"
syncpkg "git.fabledsword.com/bvandeusen/minstrel/internal/sync" syncpkg "git.fabledsword.com/bvandeusen/minstrel/internal/sync"
) )
@@ -101,12 +102,37 @@ func runDuplicateSweep(
} }
} }
if runErr == nil {
notifyNewDuplicates(finishCtx, q, sweep.StartedAt, logger)
}
logger.Info("duplicate sweep complete", logger.Info("duplicate sweep complete",
"candidates", res.Candidates, "groups", res.Groups, "proposed", res.Proposed, "candidates", res.Candidates, "groups", res.Groups, "proposed", res.Proposed,
"suppressed", res.Suppressed, "retired", res.Retired, "oversize", res.Oversize, "err", runErr) "suppressed", res.Suppressed, "retired", res.Retired, "oversize", res.Oversize, "err", runErr)
return res, runErr return res, runErr
} }
// notifyNewDuplicates tells admins when a sweep proposed a group it had not
// proposed before (M489), counting every proposal awaiting review. A sweep
// that only re-finds known groups says nothing, so reading the notice once
// is enough until something new turns up.
func notifyNewDuplicates(ctx context.Context, q *dbq.Queries, sweepStarted pgtype.Timestamptz, logger *slog.Logger) {
fresh, err := q.CountDuplicateGroupsDetectedSince(ctx, sweepStarted)
if err != nil {
logger.Warn("duplicate sweep: counting new proposals failed", "err", err)
return
}
if fresh == 0 {
return
}
pending, err := q.CountPendingDuplicateGroups(ctx)
if err != nil {
logger.Warn("duplicate sweep: counting pending proposals failed", "err", err)
return
}
notifyAdmins(ctx, notifications.KindDuplicatesFound, notifications.Payload{Count: pending})
}
func sweepDuplicates( func sweepDuplicates(
ctx context.Context, q *dbq.Queries, sweepID pgtype.UUID, cfg FingerprintSettings, pageSize int32, ctx context.Context, q *dbq.Queries, sweepID pgtype.UUID, cfg FingerprintSettings, pageSize int32,
) (DuplicateSweepResult, error) { ) (DuplicateSweepResult, error) {
+44
View File
@@ -0,0 +1,44 @@
package library
import (
"context"
"sync"
"github.com/jackc/pgx/v5/pgtype"
"git.fabledsword.com/bvandeusen/minstrel/internal/notifications"
)
// Package-level notifier for library-health notices to admins (M489): a
// scan that failed, tracks gone missing, new duplicates. Set once at startup
// beside SetEventBus, for the same reason the bus is package-level. Nil, as in
// tests that never set it, means nobody is told.
var (
notifierMu sync.RWMutex
notifier *notifications.Notifier
)
// SetNotifier wires the notifications inbox into the library package.
func SetNotifier(n *notifications.Notifier) {
notifierMu.Lock()
defer notifierMu.Unlock()
notifier = n
}
// notifyAdmins tells every admin about a library-health event. These kinds
// coalesce, so a burst is one unread notice with a running count.
func notifyAdmins(ctx context.Context, kind notifications.Kind, p notifications.Payload) {
notifierMu.RLock()
n := notifier
notifierMu.RUnlock()
n.NotifyLogged(ctx, kind, notifications.ToAdmins(pgtype.UUID{}), p.Map())
}
// notifyScanFinished tells admins a scan run ended in error. A scan cut short
// by shutdown is not a failure anyone needs telling about.
func notifyScanFinished(ctx context.Context, errMsg string) {
if errMsg == "" || ctx.Err() != nil {
return
}
notifyAdmins(ctx, notifications.KindScanFailed, notifications.Payload{Count: 1, Detail: errMsg})
}
+180
View File
@@ -0,0 +1,180 @@
package library
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"log/slog"
"path/filepath"
"testing"
"github.com/jackc/pgx/v5/pgxpool"
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
"git.fabledsword.com/bvandeusen/minstrel/internal/dbtest"
"git.fabledsword.com/bvandeusen/minstrel/internal/notifications"
)
// notifyingAdmin wires a real notifier for the test and returns an admin whose
// inbox the test reads.
func notifyingAdmin(t *testing.T, pool *pgxpool.Pool) dbq.User {
t.Helper()
admin, err := dbq.New(pool).CreateUser(context.Background(), dbq.CreateUserParams{
Username: dbtest.TestUserPrefix + "libnotify", PasswordHash: "x", ApiTokenHash: "x", IsAdmin: true,
})
if err != nil {
t.Fatalf("admin: %v", err)
}
SetNotifier(notifications.New(pool, nil, nil))
t.Cleanup(func() { SetNotifier(nil) })
return admin
}
type notice struct {
kind string
count int64
body string
}
func unreadNotices(t *testing.T, pool *pgxpool.Pool, user dbq.User) []notice {
t.Helper()
rows, err := dbq.New(pool).ListNotifications(context.Background(), dbq.ListNotificationsParams{
UserID: user.ID, PageLimit: 50,
})
if err != nil {
t.Fatalf("list notifications: %v", err)
}
var out []notice
for _, r := range rows {
if r.ReadAt.Valid {
continue
}
var p notifications.Payload
if err := json.Unmarshal(r.Payload, &p); err != nil {
t.Fatalf("payload: %v", err)
}
out = append(out, notice{kind: r.Kind, count: p.Count, body: p.Detail})
}
return out
}
func TestNotifyScanFinished_FailuresCoalesceAndShutdownIsSilent(t *testing.T) {
pool := newPool(t)
admin := notifyingAdmin(t, pool)
ctx := context.Background()
notifyScanFinished(ctx, "")
if got := unreadNotices(t, pool, admin); len(got) != 0 {
t.Fatalf("a clean scan notified: %+v", got)
}
cancelled, cancel := context.WithCancel(ctx)
cancel()
notifyScanFinished(cancelled, "library: context canceled")
if got := unreadNotices(t, pool, admin); len(got) != 0 {
t.Fatalf("a scan cut short by shutdown notified: %+v", got)
}
notifyScanFinished(ctx, "library: root /music missing")
notifyScanFinished(ctx, "library: root /music still missing")
got := unreadNotices(t, pool, admin)
if len(got) != 1 || got[0].kind != string(notifications.KindScanFailed) || got[0].count != 2 ||
got[0].body != "library: root /music still missing" {
t.Fatalf("notices = %+v, want one scan_failed counting 2 with the latest error", got)
}
}
func TestReconcileMissing_NotifiesAdminsWithARunningCount(t *testing.T) {
pool := newPool(t)
admin := notifyingAdmin(t, pool)
s := testScanner(t, populatedRoot(t))
// Two scans, each losing 2 of 10 tracks (under the mark cap).
for pass := 0; pass < 2; pass++ {
rows := make([]dbq.ListTrackPathsForReconcileRow, 0, 10)
seen := map[string]struct{}{}
for i := 0; i < 10; i++ {
p := fmt.Sprintf("/music/pass-%d-%02d.mp3", pass, i)
rows = append(rows, row(byte(pass*10+i), p, false))
if i >= 2 {
seen[p] = struct{}{}
}
}
var stats Stats
if err := s.reconcileMissing(context.Background(), &fakeReconciler{rows: rows}, seen, &stats); err != nil {
t.Fatalf("reconcile pass %d: %v", pass, err)
}
}
got := unreadNotices(t, pool, admin)
if len(got) != 1 || got[0].kind != string(notifications.KindTracksMissing) || got[0].count != 4 {
t.Fatalf("notices = %+v, want one tracks_missing counting 4", got)
}
}
func TestDuplicateSweep_NotifiesOnlyWhenSomethingNewIsProposed(t *testing.T) {
pool := newPool(t)
admin := notifyingAdmin(t, pool)
ctx := context.Background()
q := dbq.New(pool)
dir := t.TempDir()
logger := slog.New(slog.NewTextHandler(io.Discard, nil))
_, album, artist := seedTrack(t, pool, filepath.Join(dir, "seed.mp3"))
// Same audio-stream hash: an exact pair, whatever the prints say.
pair := func(name string, b byte, seed uint64) {
t.Helper()
for i := 1; i <= 2; i++ {
tr, err := q.UpsertTrack(ctx, dbq.UpsertTrackParams{
Title: name, AlbumID: album.ID, ArtistID: artist.ID, DurationMs: 200000,
FilePath: filepath.Join(dir, fmt.Sprintf("%s-%d.mp3", name, i)), FileSize: 100, FileFormat: "mp3",
})
if err != nil {
t.Fatalf("track %s: %v", name, err)
}
if err := q.UpsertTrackFingerprint(ctx, dbq.UpsertTrackFingerprintParams{
TrackID: tr.ID, AudioStreamSha256: bytes.Repeat([]byte{b}, 32),
Chromaprint: randomPrint(seed+uint64(i), printLen), FingerprintVersion: fingerprintVersion,
ChromaprintLengthSec: defaultChromaprintLengthSec,
}); err != nil {
t.Fatalf("fingerprint %s: %v", name, err)
}
}
}
sweep := func() {
t.Helper()
if _, err := runDuplicateSweep(ctx, pool, logger, DefaultFingerprintSettings, duplicateCandidatePage); err != nil {
t.Fatalf("sweep: %v", err)
}
}
markAllRead := func() {
t.Helper()
if _, err := q.MarkAllNotificationsRead(ctx, admin.ID); err != nil {
t.Fatalf("mark read: %v", err)
}
}
pair("first", 1, 100)
sweep()
got := unreadNotices(t, pool, admin)
if len(got) != 1 || got[0].kind != string(notifications.KindDuplicatesFound) || got[0].count != 1 {
t.Fatalf("after the first sweep notices = %+v, want one duplicates_found counting 1", got)
}
// Read, then swept again with nothing new: no reminder.
markAllRead()
sweep()
if got := unreadNotices(t, pool, admin); len(got) != 0 {
t.Fatalf("a sweep that found nothing new notified: %+v", got)
}
// A new pair is news, and the count is everything awaiting review.
pair("second", 2, 200)
sweep()
got = unreadNotices(t, pool, admin)
if len(got) != 1 || got[0].count != 2 {
t.Fatalf("after a new pair notices = %+v, want one counting 2", got)
}
}
+4
View File
@@ -9,6 +9,7 @@ import (
"github.com/jackc/pgx/v5/pgtype" "github.com/jackc/pgx/v5/pgtype"
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
"git.fabledsword.com/bvandeusen/minstrel/internal/notifications"
syncpkg "git.fabledsword.com/bvandeusen/minstrel/internal/sync" syncpkg "git.fabledsword.com/bvandeusen/minstrel/internal/sync"
) )
@@ -128,6 +129,9 @@ func (s *Scanner) reconcileMissing(
// this line. // this line.
s.logger.Warn("library scan: tracks marked missing (files not found)", s.logger.Warn("library scan: tracks marked missing (files not found)",
"count", n, "library_total", len(rows)) "count", n, "library_total", len(rows))
if n > 0 {
notifyAdmins(ctx, notifications.KindTracksMissing, notifications.Payload{Count: n})
}
return nil return nil
} }
+1
View File
@@ -245,6 +245,7 @@ func RunScan(
} }
logger.Info("scan run complete", "id", row.ID, "error", errMsg) logger.Info("scan run complete", "id", row.ID, "error", errMsg)
notifyScanFinished(ctx, errMsg)
publishScanEvent("scan.run_finished", row.ID, map[string]any{ publishScanEvent("scan.run_finished", row.ID, map[string]any{
"error_message": errMsg, "error_message": errMsg,
}) })