diff --git a/cmd/minstrel/main.go b/cmd/minstrel/main.go index f88a5b62..0b942d62 100644 --- a/cmd/minstrel/main.go +++ b/cmd/minstrel/main.go @@ -321,6 +321,7 @@ func run() error { lidarrReconciler := lidarrrequests.NewReconciler(pool, lidarrCfg, lidarrClientFn, logger.With("component", "lidarr"), bus) lidarrReconciler.SetNotifier(notifier) + library.SetNotifier(notifier) go lidarrReconciler.Run(ctx) // Missing-file re-acquisition (milestone #290). Turns albums whose files diff --git a/internal/api/notify_producers.go b/internal/api/notify_producers.go index b8fe22c2..e690f581 100644 --- a/internal/api/notify_producers.go +++ b/internal/api/notify_producers.go @@ -71,3 +71,16 @@ func (h *handlers) trackLabel(ctx context.Context, trackID pgtype.UUID) string { } 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()) +} diff --git a/internal/api/notify_producers_test.go b/internal/api/notify_producers_test.go index d90db858..d5f7d0e7 100644 --- a/internal/api/notify_producers_test.go +++ b/internal/api/notify_producers_test.go @@ -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)) +} diff --git a/internal/api/playback_errors.go b/internal/api/playback_errors.go index fb7da802..ccc358c7 100644 --- a/internal/api/playback_errors.go +++ b/internal/api/playback_errors.go @@ -130,6 +130,7 @@ func (h *handlers) handleReportPlaybackError(w http.ResponseWriter, r *http.Requ writeErr(w, apierror.Internal(err)) return } + h.notifyPlaybackErrors(r.Context()) writeJSON(w, http.StatusCreated, map[string]string{"id": uuidToString(row.ID)}) } diff --git a/internal/db/dbq/duplicates.sql.go b/internal/db/dbq/duplicates.sql.go index c3dab489..beb0161e 100644 --- a/internal/db/dbq/duplicates.sql.go +++ b/internal/db/dbq/duplicates.sql.go @@ -27,6 +27,24 @@ func (q *Queries) AddDuplicateGroupMember(ctx context.Context, arg AddDuplicateG 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 SELECT count(*)::bigint FROM duplicate_groups g diff --git a/internal/db/dbq/playback_errors.sql.go b/internal/db/dbq/playback_errors.sql.go index 607441b9..53988c4e 100644 --- a/internal/db/dbq/playback_errors.sql.go +++ b/internal/db/dbq/playback_errors.sql.go @@ -11,6 +11,17 @@ import ( "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 INSERT INTO playback_errors (track_id, user_id, client_id, kind, detail) VALUES ($1, $2, $3, $4, $5) diff --git a/internal/db/queries/duplicates.sql b/internal/db/queries/duplicates.sql index 9150ccc5..0aada774 100644 --- a/internal/db/queries/duplicates.sql +++ b/internal/db/queries/duplicates.sql @@ -154,3 +154,13 @@ SELECT p.id AS group_id, UPDATE duplicate_groups SET status = 'dismissed', resolved_at = now() 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); diff --git a/internal/db/queries/playback_errors.sql b/internal/db/queries/playback_errors.sql index c89da41f..823a76ee 100644 --- a/internal/db/queries/playback_errors.sql +++ b/internal/db/queries/playback_errors.sql @@ -61,3 +61,6 @@ UPDATE playback_errors resolved_by = $2, resolution = $3 WHERE track_id = $1 AND resolved_at IS NULL; + +-- name: CountUnresolvedPlaybackErrors :one +SELECT count(*)::bigint FROM playback_errors WHERE resolved_at IS NULL; diff --git a/internal/library/duplicate_sweep.go b/internal/library/duplicate_sweep.go index 34dd2cb8..46e6cba7 100644 --- a/internal/library/duplicate_sweep.go +++ b/internal/library/duplicate_sweep.go @@ -14,6 +14,7 @@ import ( "github.com/jackc/pgx/v5/pgxpool" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" 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", "candidates", res.Candidates, "groups", res.Groups, "proposed", res.Proposed, "suppressed", res.Suppressed, "retired", res.Retired, "oversize", res.Oversize, "err", 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( ctx context.Context, q *dbq.Queries, sweepID pgtype.UUID, cfg FingerprintSettings, pageSize int32, ) (DuplicateSweepResult, error) { diff --git a/internal/library/notify.go b/internal/library/notify.go new file mode 100644 index 00000000..01b4f361 --- /dev/null +++ b/internal/library/notify.go @@ -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}) +} diff --git a/internal/library/notify_test.go b/internal/library/notify_test.go new file mode 100644 index 00000000..eb7b4ba9 --- /dev/null +++ b/internal/library/notify_test.go @@ -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) + } +} diff --git a/internal/library/reconcile.go b/internal/library/reconcile.go index 534b42e8..1f13474b 100644 --- a/internal/library/reconcile.go +++ b/internal/library/reconcile.go @@ -9,6 +9,7 @@ import ( "github.com/jackc/pgx/v5/pgtype" "git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq" + "git.fabledsword.com/bvandeusen/minstrel/internal/notifications" syncpkg "git.fabledsword.com/bvandeusen/minstrel/internal/sync" ) @@ -128,6 +129,9 @@ func (s *Scanner) reconcileMissing( // this line. s.logger.Warn("library scan: tracks marked missing (files not found)", "count", n, "library_total", len(rows)) + if n > 0 { + notifyAdmins(ctx, notifications.KindTracksMissing, notifications.Payload{Count: n}) + } return nil } diff --git a/internal/library/scanrun.go b/internal/library/scanrun.go index e907eaab..f25081f8 100644 --- a/internal/library/scanrun.go +++ b/internal/library/scanrun.go @@ -245,6 +245,7 @@ func RunScan( } logger.Info("scan run complete", "id", row.ID, "error", errMsg) + notifyScanFinished(ctx, errMsg) publishScanEvent("scan.run_finished", row.ID, map[string]any{ "error_message": errMsg, })