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>
425 lines
15 KiB
Go
425 lines
15 KiB
Go
package library
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgtype"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
|
|
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
|
|
"git.fabledsword.com/bvandeusen/minstrel/internal/notifications"
|
|
syncpkg "git.fabledsword.com/bvandeusen/minstrel/internal/sync"
|
|
)
|
|
|
|
// Duplicate sweep (M400 #3910).
|
|
//
|
|
// Reads fingerprints, runs them through the matcher and records what it proposes
|
|
// in duplicate_groups. It never merges or deletes anything the operator has not
|
|
// asked for: a group is a proposal, reviewed in the admin report (#3912).
|
|
//
|
|
// Blocked on duration, deliberately not on title: the #3885 pair are titled
|
|
// "WWW" and "WWW (instrumental)", so a title block would have missed the case
|
|
// that started the milestone. Candidates stream in (duration_ms, id) order and
|
|
// the grouper holds only a few seconds of durations at a time.
|
|
|
|
// duplicateCandidatePage is how many candidates one query returns. Each row
|
|
// carries a ~4 KB fingerprint, so a page is about 2 MB.
|
|
const duplicateCandidatePage = 500
|
|
|
|
// duplicateSweepTick is how often the worker checks whether a sweep is due. The
|
|
// operator's sweep interval (#3913) is the least time between sweeps; the tick
|
|
// only bounds how late past it one starts. With nothing due, a tick is two cheap
|
|
// aggregate queries.
|
|
const duplicateSweepTick = time.Hour
|
|
|
|
// sweepIntervalSlack absorbs the moment between a tick and the sweep it starts
|
|
// stamping started_at. Without it a one-hour interval checked on a one-hour tick
|
|
// would find the last sweep a moment under an hour old, and skip every other tick.
|
|
const sweepIntervalSlack = 5 * time.Minute
|
|
|
|
// staleDuplicateSweepThreshold is the age past which an in-flight sweep is
|
|
// assumed dead — a crash mid-sweep leaves finished_at NULL for ever — and another
|
|
// may start. Twice the library scan's threshold, because a sweep compares
|
|
// fingerprints across the whole library and can legitimately run long on a big
|
|
// one.
|
|
const staleDuplicateSweepThreshold = 2 * time.Hour
|
|
|
|
// duplicateSweepFinishTimeout bounds recording that a sweep ended. It runs on a
|
|
// context detached from the sweep's own, so a sweep cancelled at shutdown still
|
|
// closes its row rather than leaving it in flight until the reaper.
|
|
const duplicateSweepFinishTimeout = 10 * time.Second
|
|
|
|
// DuplicateSweepResult tallies one sweep.
|
|
type DuplicateSweepResult struct {
|
|
Candidates int // tracks with a chromaprint that were streamed
|
|
Groups int // groups the matcher found
|
|
Proposed int // written as pending, new or refreshed
|
|
Suppressed int // not proposed: already dismissed or resolved by the operator
|
|
Retired int // pending proposals this sweep did not find again, removed
|
|
Oversize int // acoustic clusters too large to propose
|
|
}
|
|
|
|
// RunDuplicateSweep runs one sweep and records it in duplicate_sweeps. cfg is a
|
|
// snapshot: one sweep applies one threshold and one length throughout.
|
|
func RunDuplicateSweep(
|
|
ctx context.Context, pool *pgxpool.Pool, logger *slog.Logger, cfg FingerprintSettings,
|
|
) (DuplicateSweepResult, error) {
|
|
return runDuplicateSweep(ctx, pool, logger, cfg, duplicateCandidatePage)
|
|
}
|
|
|
|
func runDuplicateSweep(
|
|
ctx context.Context, pool *pgxpool.Pool, logger *slog.Logger, cfg FingerprintSettings, pageSize int32,
|
|
) (DuplicateSweepResult, error) {
|
|
q := dbq.New(pool)
|
|
sweep, err := q.StartDuplicateSweep(ctx)
|
|
if err != nil {
|
|
return DuplicateSweepResult{}, fmt.Errorf("start duplicate sweep: %w", err)
|
|
}
|
|
|
|
res, runErr := sweepDuplicates(ctx, q, sweep.ID, cfg, pageSize)
|
|
|
|
finishCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), duplicateSweepFinishTimeout)
|
|
defer cancel()
|
|
errMsg := ""
|
|
if runErr != nil {
|
|
errMsg = runErr.Error()
|
|
}
|
|
candidates, groups, oversize := int32(res.Candidates), int32(res.Groups), int32(res.Oversize)
|
|
if ferr := q.FinishDuplicateSweep(finishCtx, dbq.FinishDuplicateSweepParams{
|
|
ID: sweep.ID, Candidates: &candidates, GroupsFound: &groups, OversizeClusters: &oversize,
|
|
ErrorMessage: errMsg,
|
|
}); ferr != nil {
|
|
logger.Error("duplicate sweep: recording the end of the sweep failed", "err", ferr)
|
|
if runErr == nil {
|
|
runErr = fmt.Errorf("finish duplicate sweep: %w", ferr)
|
|
}
|
|
}
|
|
|
|
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) {
|
|
var res DuplicateSweepResult
|
|
|
|
// Exact tier, library-wide, in one query.
|
|
exactRows, err := q.ListExactDuplicateHashes(ctx, fingerprintVersion)
|
|
if err != nil {
|
|
return res, fmt.Errorf("list exact duplicates: %w", err)
|
|
}
|
|
exactMembers := make([][]string, len(exactRows))
|
|
exactOf := map[string]int{}
|
|
for i, row := range exactRows {
|
|
exactMembers[i] = formatUUIDs(row.TrackIds)
|
|
for _, id := range exactMembers[i] {
|
|
exactOf[id] = i
|
|
}
|
|
}
|
|
exactSeen := make([]bool, len(exactRows))
|
|
|
|
// Acoustic tier, streamed in duration order. The first member of an exact
|
|
// group the stream meets stands in for the whole group; the rest are skipped.
|
|
// Only prints at the current length are streamed: a print at another length
|
|
// cannot be compared, and is waiting on the backfill to be re-derived.
|
|
grouper := newStreamGrouper(cfg.AcousticMaxBitErrorRate)
|
|
params := dbq.ListDuplicateCandidatesParams{
|
|
CurrentVersion: fingerprintVersion,
|
|
ChromaprintLengthSec: cfg.ChromaprintLengthSec,
|
|
// Durations are never negative, and the all-zero uuid sorts first: every
|
|
// row is after this cursor. Valid must be true, or "> NULL" matches nothing.
|
|
AfterDurationMs: -1,
|
|
AfterID: pgtype.UUID{Valid: true},
|
|
PageLimit: pageSize,
|
|
}
|
|
for {
|
|
if err := ctx.Err(); err != nil {
|
|
return res, err
|
|
}
|
|
rows, err := q.ListDuplicateCandidates(ctx, params)
|
|
if err != nil {
|
|
return res, fmt.Errorf("list duplicate candidates: %w", err)
|
|
}
|
|
for _, row := range rows {
|
|
res.Candidates++
|
|
id := syncpkg.FormatUUID(row.ID)
|
|
unit := &groupUnit{ids: []string{id}, durationMs: row.DurationMs, sortKey: id}
|
|
if gi, ok := exactOf[id]; ok {
|
|
if exactSeen[gi] {
|
|
continue
|
|
}
|
|
exactSeen[gi] = true
|
|
unit.ids, unit.exact = exactMembers[gi], true
|
|
}
|
|
unit.print = preparePrint(row.Chromaprint)
|
|
grouper.push(unit)
|
|
}
|
|
if int32(len(rows)) < pageSize {
|
|
break
|
|
}
|
|
last := rows[len(rows)-1]
|
|
params.AfterDurationMs, params.AfterID = last.DurationMs, last.ID
|
|
}
|
|
|
|
found := grouper.finish()
|
|
// Exact groups none of whose members has a chromaprint never reached the
|
|
// stream. Identical bytes need no acoustic evidence.
|
|
for gi, seen := range exactSeen {
|
|
if !seen {
|
|
found.Groups = append(found.Groups, duplicateGroup{Tier: tierExact, Members: exactMembers[gi]})
|
|
}
|
|
}
|
|
sortGroups(found.Groups)
|
|
res.Groups, res.Oversize = len(found.Groups), found.OversizeClusters
|
|
|
|
dismissed, err := q.ListDismissedDuplicateMemberSets(ctx)
|
|
if err != nil {
|
|
return res, fmt.Errorf("list dismissed duplicate groups: %w", err)
|
|
}
|
|
dismissedSets := make([]map[string]struct{}, 0, len(dismissed))
|
|
for _, d := range dismissed {
|
|
set := map[string]struct{}{}
|
|
for _, id := range formatUUIDs(d.TrackIds) {
|
|
set[id] = struct{}{}
|
|
}
|
|
dismissedSets = append(dismissedSets, set)
|
|
}
|
|
|
|
for _, group := range found.Groups {
|
|
if coveredByDismissal(group.Members, dismissedSets) {
|
|
res.Suppressed++
|
|
continue
|
|
}
|
|
up := dbq.UpsertDuplicateGroupParams{
|
|
MemberKey: strings.Join(group.Members, ","),
|
|
Tier: string(group.Tier),
|
|
SweepID: sweepID,
|
|
}
|
|
if group.Tier == tierAcoustic {
|
|
worst := float32(group.WorstBitErrorRate)
|
|
up.WorstBitErrorRate = &worst
|
|
}
|
|
groupID, err := q.UpsertDuplicateGroup(ctx, up)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
// This exact member set was already dismissed or merged.
|
|
res.Suppressed++
|
|
continue
|
|
}
|
|
if err != nil {
|
|
return res, fmt.Errorf("upsert duplicate group: %w", err)
|
|
}
|
|
for _, id := range group.Members {
|
|
var trackID pgtype.UUID
|
|
if err := trackID.Scan(id); err != nil {
|
|
return res, fmt.Errorf("parse track id %q: %w", id, err)
|
|
}
|
|
if err := q.AddDuplicateGroupMember(ctx, dbq.AddDuplicateGroupMemberParams{
|
|
GroupID: groupID, TrackID: trackID,
|
|
}); err != nil {
|
|
return res, fmt.Errorf("add duplicate group member: %w", err)
|
|
}
|
|
}
|
|
res.Proposed++
|
|
}
|
|
|
|
// Only after a complete sweep: a sweep that failed partway has no basis for
|
|
// concluding that anything it did not reach has gone.
|
|
retired, err := q.DeleteStalePendingDuplicateGroups(ctx, sweepID)
|
|
if err != nil {
|
|
return res, fmt.Errorf("retire stale duplicate groups: %w", err)
|
|
}
|
|
res.Retired = int(retired)
|
|
return res, nil
|
|
}
|
|
|
|
// coveredByDismissal reports whether every member of a proposal sat together in
|
|
// one group the operator dismissed. A subset counts: dismissing {A, B, C} said
|
|
// none of them are copies of each other, so proposing {A, B} again would be
|
|
// asking the same question twice. A superset does not count: a new copy joining
|
|
// is new evidence, and worth asking about.
|
|
func coveredByDismissal(members []string, dismissed []map[string]struct{}) bool {
|
|
for _, set := range dismissed {
|
|
covered := true
|
|
for _, id := range members {
|
|
if _, ok := set[id]; !ok {
|
|
covered = false
|
|
break
|
|
}
|
|
}
|
|
if covered {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
func formatUUIDs(ids []pgtype.UUID) []string {
|
|
out := make([]string, len(ids))
|
|
for i, id := range ids {
|
|
out[i] = syncpkg.FormatUUID(id)
|
|
}
|
|
sort.Strings(out)
|
|
return out
|
|
}
|
|
|
|
// TryStartDuplicateSweep starts a sweep in the background unless one is already
|
|
// running, reaping a sweep that has been in flight past
|
|
// staleDuplicateSweepThreshold. Mirrors TryStartScan. The sweep runs on ctx, so
|
|
// a caller answering an HTTP request must pass a context that outlives it.
|
|
func TryStartDuplicateSweep(
|
|
ctx context.Context, pool *pgxpool.Pool, logger *slog.Logger, cfg FingerprintSettings,
|
|
) (bool, error) {
|
|
q := dbq.New(pool)
|
|
row, err := q.GetInFlightDuplicateSweep(ctx)
|
|
switch {
|
|
case err == nil:
|
|
age := time.Since(row.StartedAt.Time)
|
|
if age <= staleDuplicateSweepThreshold {
|
|
return false, nil
|
|
}
|
|
logger.Warn("reaping stale duplicate sweep", "id", syncpkg.FormatUUID(row.ID), "age", age)
|
|
if ferr := q.FinishDuplicateSweep(ctx, dbq.FinishDuplicateSweepParams{
|
|
ID: row.ID, ErrorMessage: "reaped (stale)",
|
|
}); ferr != nil {
|
|
return false, fmt.Errorf("reap stale duplicate sweep: %w", ferr)
|
|
}
|
|
case !errors.Is(err, pgx.ErrNoRows):
|
|
return false, fmt.Errorf("duplicate sweep in-flight check: %w", err)
|
|
}
|
|
|
|
go func() {
|
|
if _, err := RunDuplicateSweep(ctx, pool, logger, cfg); err != nil {
|
|
logger.Warn("duplicate sweep failed", "err", err)
|
|
}
|
|
}()
|
|
return true, nil
|
|
}
|
|
|
|
// DuplicateSweepWorker sweeps whenever its input has changed, at most once per
|
|
// the operator's sweep interval.
|
|
type DuplicateSweepWorker struct {
|
|
pool *pgxpool.Pool
|
|
logger *slog.Logger
|
|
settings *FingerprintSettingsService
|
|
tick time.Duration
|
|
}
|
|
|
|
// NewDuplicateSweepWorker builds a worker with the production cadence. settings
|
|
// is shared with the admin API; nil runs on defaults.
|
|
func NewDuplicateSweepWorker(
|
|
pool *pgxpool.Pool, logger *slog.Logger, settings *FingerprintSettingsService,
|
|
) *DuplicateSweepWorker {
|
|
return &DuplicateSweepWorker{pool: pool, logger: logger, settings: settings, tick: duplicateSweepTick}
|
|
}
|
|
|
|
// Run blocks until ctx is cancelled, checking once at start and then each tick.
|
|
func (w *DuplicateSweepWorker) Run(ctx context.Context) {
|
|
w.tickOnce(ctx)
|
|
t := time.NewTicker(w.tick)
|
|
defer t.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-t.C:
|
|
w.tickOnce(ctx)
|
|
}
|
|
}
|
|
}
|
|
|
|
// tickOnce contains one check so nothing it does can stop the next tick (rule 157).
|
|
func (w *DuplicateSweepWorker) tickOnce(ctx context.Context) {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
w.logger.Error("duplicate sweep: tick panicked", "panic", r)
|
|
}
|
|
}()
|
|
cfg := w.settings.Get()
|
|
due, err := duplicateSweepDue(ctx, dbq.New(w.pool), cfg, time.Now())
|
|
if err != nil {
|
|
if ctx.Err() == nil {
|
|
w.logger.Warn("duplicate sweep: due check failed", "err", err)
|
|
}
|
|
return
|
|
}
|
|
if !due {
|
|
return
|
|
}
|
|
if _, err := TryStartDuplicateSweep(ctx, w.pool, w.logger, cfg); err != nil {
|
|
w.logger.Warn("duplicate sweep: start failed", "err", err)
|
|
}
|
|
}
|
|
|
|
// duplicateSweepDue reads what sweepIsDue decides on.
|
|
func duplicateSweepDue(ctx context.Context, q *dbq.Queries, cfg FingerprintSettings, now time.Time) (bool, error) {
|
|
latest, err := q.GetLatestFingerprintComputedAt(ctx)
|
|
if err != nil {
|
|
return false, fmt.Errorf("latest fingerprint: %w", err)
|
|
}
|
|
var lastStart pgtype.Timestamptz
|
|
last, err := q.GetLatestDuplicateSweep(ctx)
|
|
switch {
|
|
case err == nil:
|
|
lastStart = last.StartedAt
|
|
case !errors.Is(err, pgx.ErrNoRows):
|
|
return false, fmt.Errorf("latest duplicate sweep: %w", err)
|
|
}
|
|
return sweepIsDue(latest, lastStart, cfg, now), nil
|
|
}
|
|
|
|
// sweepIsDue reports whether a sweep should start: something it reads has
|
|
// changed since the last sweep started, and the operator's interval has passed.
|
|
//
|
|
// Two things can change its answer. Fingerprints are its input, so any written
|
|
// after the last sweep began count; while the backfill runs that is true every
|
|
// tick, which is what the interval is for. And a settings save counts, because a
|
|
// new threshold or length changes what the same fingerprints group into.
|
|
func sweepIsDue(latestFingerprint, lastSweepStart pgtype.Timestamptz, cfg FingerprintSettings, now time.Time) bool {
|
|
if !latestFingerprint.Valid {
|
|
return false // nothing fingerprinted yet
|
|
}
|
|
if !lastSweepStart.Valid {
|
|
return true // never swept
|
|
}
|
|
interval := time.Duration(cfg.SweepIntervalHours) * time.Hour
|
|
if now.Sub(lastSweepStart.Time) < interval-sweepIntervalSlack {
|
|
return false
|
|
}
|
|
return latestFingerprint.Time.After(lastSweepStart.Time) || cfg.UpdatedAt.After(lastSweepStart.Time)
|
|
}
|