Files
minstrel/internal/db/dbq/duplicates.sql.go
T
bvandeusenandClaude Opus 5 077ae61235
test-go / test (push) Failing after 44s
test-web / test (push) Successful in 49s
test-go / integration (push) Failing after 2m42s
release / Build + push container image (push) Canceled after 0s
release / Verify release artifacts (tag releases only) (push) Canceled after 0s
release / Build signed APK (releases and dev) (push) Canceled after 4m8s
feat(admin): fingerprinting settings — on/off, length, match threshold, concurrency, sweep interval (M400 #3913)
Rule 25: the fingerprinting knobs move out of source into a DB-backed
singleton (migration 0061), edited from a card on the Duplicates page and
shared live with the scanner, the backfill and the duplicate sweep through
one service instance, so a save needs no restart.

The length is the knob that can silently break the library: prints taken
at two lengths never match. Each track_fingerprints row now records the
length it was taken at, and every reader filters on the current one — the
backfill treats another length as stale, the gauge counts it pending, the
sweep never streams it. Equivalent to a version bump, except that setting
the length back makes rows not yet redone current again. The card warns
before a length change re-fingerprints the library.

Off stops every decode: the scan takes only the stream hash (a demux, and
what recognises a moved file) and stores nothing, dropping a changed file's
stale row; the backfill idles. A save also makes a sweep due, since a new
threshold or length changes what the same prints group into, and the sweep
interval gains slack so an hourly interval on an hourly tick doesn't skip
every other tick.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01SQ31KQpYbStyK5y58UmPLH
2026-09-11 17:53:55 -04:00

458 lines
14 KiB
Go

// Code generated by sqlc. DO NOT EDIT.
// versions:
// sqlc v1.31.1
// source: duplicates.sql
package dbq
import (
"context"
"github.com/jackc/pgx/v5/pgtype"
)
const addDuplicateGroupMember = `-- name: AddDuplicateGroupMember :exec
INSERT INTO duplicate_group_members (group_id, track_id)
VALUES ($1, $2)
ON CONFLICT DO NOTHING
`
type AddDuplicateGroupMemberParams struct {
GroupID pgtype.UUID
TrackID pgtype.UUID
}
func (q *Queries) AddDuplicateGroupMember(ctx context.Context, arg AddDuplicateGroupMemberParams) error {
_, err := q.db.Exec(ctx, addDuplicateGroupMember, arg.GroupID, arg.TrackID)
return err
}
const countPendingDuplicateGroups = `-- name: CountPendingDuplicateGroups :one
SELECT count(*)::bigint
FROM duplicate_groups g
WHERE g.status = 'pending'
AND (SELECT count(*) FROM duplicate_group_members m WHERE m.group_id = g.id) >= 2
`
// Proposals awaiting review. A group left with one member — its other tracks
// deleted since the sweep — is no proposal at all and is not counted; the next
// sweep retires it.
func (q *Queries) CountPendingDuplicateGroups(ctx context.Context) (int64, error) {
row := q.db.QueryRow(ctx, countPendingDuplicateGroups)
var column_1 int64
err := row.Scan(&column_1)
return column_1, err
}
const deleteStalePendingDuplicateGroups = `-- name: DeleteStalePendingDuplicateGroups :execrows
DELETE FROM duplicate_groups g
WHERE g.status = 'pending'
AND g.last_seen_sweep_id IS DISTINCT FROM $1
AND (g.last_seen_sweep_id IS NULL
OR (SELECT s.started_at FROM duplicate_sweeps s WHERE s.id = g.last_seen_sweep_id)
< (SELECT s.started_at FROM duplicate_sweeps s WHERE s.id = $1))
`
// A pending proposal this sweep did not find again no longer describes the
// library: a member was re-fingerprinted, merged away or went missing. Dismissed
// groups are kept regardless — they are the memory of a decision.
//
// Only proposals last confirmed by an EARLIER sweep go. Should two sweeps ever
// overlap (a manual trigger racing the worker), neither may delete what the other
// has just found.
func (q *Queries) DeleteStalePendingDuplicateGroups(ctx context.Context, sweepID pgtype.UUID) (int64, error) {
result, err := q.db.Exec(ctx, deleteStalePendingDuplicateGroups, sweepID)
if err != nil {
return 0, err
}
return result.RowsAffected(), nil
}
const dismissDuplicateGroup = `-- name: DismissDuplicateGroup :execrows
UPDATE duplicate_groups
SET status = 'dismissed', resolved_at = now()
WHERE id = $1 AND status = 'pending'
`
// "These are not duplicates." Only a pending group can be dismissed; zero rows
// means it was already resolved or no longer exists.
func (q *Queries) DismissDuplicateGroup(ctx context.Context, id pgtype.UUID) (int64, error) {
result, err := q.db.Exec(ctx, dismissDuplicateGroup, id)
if err != nil {
return 0, err
}
return result.RowsAffected(), nil
}
const finishDuplicateSweep = `-- name: FinishDuplicateSweep :exec
UPDATE duplicate_sweeps
SET finished_at = now(),
candidates = $1,
groups_found = $2,
oversize_clusters = $3,
error_message = NULLIF($4::text, '')
WHERE id = $5
`
type FinishDuplicateSweepParams struct {
Candidates *int32
GroupsFound *int32
OversizeClusters *int32
ErrorMessage string
ID pgtype.UUID
}
func (q *Queries) FinishDuplicateSweep(ctx context.Context, arg FinishDuplicateSweepParams) error {
_, err := q.db.Exec(ctx, finishDuplicateSweep,
arg.Candidates,
arg.GroupsFound,
arg.OversizeClusters,
arg.ErrorMessage,
arg.ID,
)
return err
}
const getInFlightDuplicateSweep = `-- name: GetInFlightDuplicateSweep :one
SELECT id, started_at
FROM duplicate_sweeps
WHERE finished_at IS NULL
ORDER BY started_at DESC
LIMIT 1
`
type GetInFlightDuplicateSweepRow struct {
ID pgtype.UUID
StartedAt pgtype.Timestamptz
}
// The guard against two sweeps at once: "in flight" is finished_at IS NULL.
func (q *Queries) GetInFlightDuplicateSweep(ctx context.Context) (GetInFlightDuplicateSweepRow, error) {
row := q.db.QueryRow(ctx, getInFlightDuplicateSweep)
var i GetInFlightDuplicateSweepRow
err := row.Scan(&i.ID, &i.StartedAt)
return i, err
}
const getLatestDuplicateSweep = `-- name: GetLatestDuplicateSweep :one
SELECT id, started_at, finished_at, candidates, groups_found, oversize_clusters, error_message
FROM duplicate_sweeps
ORDER BY started_at DESC
LIMIT 1
`
func (q *Queries) GetLatestDuplicateSweep(ctx context.Context) (DuplicateSweep, error) {
row := q.db.QueryRow(ctx, getLatestDuplicateSweep)
var i DuplicateSweep
err := row.Scan(
&i.ID,
&i.StartedAt,
&i.FinishedAt,
&i.Candidates,
&i.GroupsFound,
&i.OversizeClusters,
&i.ErrorMessage,
)
return i, err
}
const getLatestFingerprintComputedAt = `-- name: GetLatestFingerprintComputedAt :one
SELECT max(computed_at)::timestamptz AS latest FROM track_fingerprints
`
// Whether a sweep has anything new to look at: fingerprints written since the
// last sweep started.
func (q *Queries) GetLatestFingerprintComputedAt(ctx context.Context) (pgtype.Timestamptz, error) {
row := q.db.QueryRow(ctx, getLatestFingerprintComputedAt)
var latest pgtype.Timestamptz
err := row.Scan(&latest)
return latest, err
}
const listDismissedDuplicateMemberSets = `-- name: ListDismissedDuplicateMemberSets :many
SELECT g.id, array_agg(m.track_id ORDER BY m.track_id)::uuid[] AS track_ids
FROM duplicate_groups g
JOIN duplicate_group_members m ON m.group_id = g.id
WHERE g.status = 'dismissed'
GROUP BY g.id
`
type ListDismissedDuplicateMemberSetsRow struct {
ID pgtype.UUID
TrackIds []pgtype.UUID
}
// What the operator has already said are not duplicates. A new proposal whose
// every member sat together in one of these is not proposed again.
func (q *Queries) ListDismissedDuplicateMemberSets(ctx context.Context) ([]ListDismissedDuplicateMemberSetsRow, error) {
rows, err := q.db.Query(ctx, listDismissedDuplicateMemberSets)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListDismissedDuplicateMemberSetsRow
for rows.Next() {
var i ListDismissedDuplicateMemberSetsRow
if err := rows.Scan(&i.ID, &i.TrackIds); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listDuplicateCandidates = `-- name: ListDuplicateCandidates :many
SELECT t.id, t.duration_ms, f.audio_stream_sha256, f.chromaprint
FROM tracks t
JOIN track_fingerprints f ON f.track_id = t.id
WHERE t.missing_since IS NULL
AND f.fingerprint_version >= $1
AND f.chromaprint IS NOT NULL
-- Only chromaprints taken at the current length: prints at two lengths are not
-- comparable, and after a length change the backfill is still re-deriving the
-- rest (#3913).
AND f.chromaprint_length_sec = $2
AND (t.duration_ms, t.id) > ($3::integer, $4::uuid)
ORDER BY t.duration_ms, t.id
LIMIT $5
`
type ListDuplicateCandidatesParams struct {
CurrentVersion int16
ChromaprintLengthSec int32
AfterDurationMs int32
AfterID pgtype.UUID
PageLimit int32
}
type ListDuplicateCandidatesRow struct {
ID pgtype.UUID
DurationMs int32
AudioStreamSha256 []byte
Chromaprint []int32
}
// The acoustic tier's input, one page at a time in (duration_ms, id) order so the
// sweep holds only a sliding window of durations. Tracks without a chromaprint
// cannot be compared acoustically and are left out; any exact duplicates among
// them come from ListExactDuplicateHashes.
func (q *Queries) ListDuplicateCandidates(ctx context.Context, arg ListDuplicateCandidatesParams) ([]ListDuplicateCandidatesRow, error) {
rows, err := q.db.Query(ctx, listDuplicateCandidates,
arg.CurrentVersion,
arg.ChromaprintLengthSec,
arg.AfterDurationMs,
arg.AfterID,
arg.PageLimit,
)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListDuplicateCandidatesRow
for rows.Next() {
var i ListDuplicateCandidatesRow
if err := rows.Scan(
&i.ID,
&i.DurationMs,
&i.AudioStreamSha256,
&i.Chromaprint,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listExactDuplicateHashes = `-- name: ListExactDuplicateHashes :many
SELECT f.audio_stream_sha256,
array_agg(t.id ORDER BY t.id)::uuid[] AS track_ids
FROM track_fingerprints f
JOIN tracks t ON t.id = f.track_id
WHERE t.missing_since IS NULL
AND f.fingerprint_version >= $1
AND f.audio_stream_sha256 IS NOT NULL
GROUP BY f.audio_stream_sha256
HAVING count(*) > 1
`
type ListExactDuplicateHashesRow struct {
AudioStreamSha256 []byte
TrackIds []pgtype.UUID
}
// The exact tier, library-wide in one pass: identical encoded audio shared by
// more than one present track.
func (q *Queries) ListExactDuplicateHashes(ctx context.Context, currentVersion int16) ([]ListExactDuplicateHashesRow, error) {
rows, err := q.db.Query(ctx, listExactDuplicateHashes, currentVersion)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListExactDuplicateHashesRow
for rows.Next() {
var i ListExactDuplicateHashesRow
if err := rows.Scan(&i.AudioStreamSha256, &i.TrackIds); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listPendingDuplicateGroupMembers = `-- name: ListPendingDuplicateGroupMembers :many
WITH page AS (
SELECT g.id, g.tier, g.worst_bit_error_rate, g.detected_at
FROM duplicate_groups g
WHERE g.status = 'pending'
AND (SELECT count(*) FROM duplicate_group_members m WHERE m.group_id = g.id) >= 2
ORDER BY g.detected_at DESC, g.id
LIMIT $2 OFFSET $1
)
SELECT p.id AS group_id,
p.tier,
p.worst_bit_error_rate,
p.detected_at,
t.id AS track_id,
t.title,
artists.name AS artist_name,
albums.id AS album_id,
albums.title AS album_title,
t.file_path,
t.file_format,
t.file_size,
t.duration_ms,
t.added_at,
(SELECT count(*) FROM general_likes l WHERE l.track_id = t.id)::bigint AS like_count,
(SELECT count(*) FROM play_events e WHERE e.track_id = t.id)::bigint AS play_count
FROM page p
JOIN duplicate_group_members m ON m.group_id = p.id
JOIN tracks t ON t.id = m.track_id
JOIN albums ON albums.id = t.album_id
JOIN artists ON artists.id = t.artist_id
ORDER BY p.detected_at DESC, p.id, t.id
`
type ListPendingDuplicateGroupMembersParams struct {
PageOffset int32
PageLimit int32
}
type ListPendingDuplicateGroupMembersRow struct {
GroupID pgtype.UUID
Tier string
WorstBitErrorRate *float32
DetectedAt pgtype.Timestamptz
TrackID pgtype.UUID
Title string
ArtistName string
AlbumID pgtype.UUID
AlbumTitle string
FilePath string
FileFormat string
FileSize int64
DurationMs int32
AddedAt pgtype.Timestamptz
LikeCount int64
PlayCount int64
}
// One page of proposals, newest first, flattened to one row per member so the
// handler folds them without a query per group. What each copy carries — likes
// and plays from every user — is here because it is what the operator weighs
// when deciding which copy to keep.
func (q *Queries) ListPendingDuplicateGroupMembers(ctx context.Context, arg ListPendingDuplicateGroupMembersParams) ([]ListPendingDuplicateGroupMembersRow, error) {
rows, err := q.db.Query(ctx, listPendingDuplicateGroupMembers, arg.PageOffset, arg.PageLimit)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListPendingDuplicateGroupMembersRow
for rows.Next() {
var i ListPendingDuplicateGroupMembersRow
if err := rows.Scan(
&i.GroupID,
&i.Tier,
&i.WorstBitErrorRate,
&i.DetectedAt,
&i.TrackID,
&i.Title,
&i.ArtistName,
&i.AlbumID,
&i.AlbumTitle,
&i.FilePath,
&i.FileFormat,
&i.FileSize,
&i.DurationMs,
&i.AddedAt,
&i.LikeCount,
&i.PlayCount,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const startDuplicateSweep = `-- name: StartDuplicateSweep :one
INSERT INTO duplicate_sweeps DEFAULT VALUES RETURNING id, started_at
`
type StartDuplicateSweepRow struct {
ID pgtype.UUID
StartedAt pgtype.Timestamptz
}
func (q *Queries) StartDuplicateSweep(ctx context.Context) (StartDuplicateSweepRow, error) {
row := q.db.QueryRow(ctx, startDuplicateSweep)
var i StartDuplicateSweepRow
err := row.Scan(&i.ID, &i.StartedAt)
return i, err
}
const upsertDuplicateGroup = `-- name: UpsertDuplicateGroup :one
INSERT INTO duplicate_groups (member_key, tier, worst_bit_error_rate, last_seen_sweep_id)
VALUES ($1, $2, $3, $4)
ON CONFLICT (member_key) DO UPDATE
SET tier = EXCLUDED.tier,
worst_bit_error_rate = EXCLUDED.worst_bit_error_rate,
last_seen_sweep_id = EXCLUDED.last_seen_sweep_id
WHERE duplicate_groups.status = 'pending'
RETURNING id
`
type UpsertDuplicateGroupParams struct {
MemberKey string
Tier string
WorstBitErrorRate *float32
SweepID pgtype.UUID
}
// Proposes a group, or refreshes one already pending. A group already dismissed
// or merged is left exactly as it is: the WHERE on the update makes the conflict
// a no-op, and the caller sees no row.
func (q *Queries) UpsertDuplicateGroup(ctx context.Context, arg UpsertDuplicateGroupParams) (pgtype.UUID, error) {
row := q.db.QueryRow(ctx, upsertDuplicateGroup,
arg.MemberKey,
arg.Tier,
arg.WorstBitErrorRate,
arg.SweepID,
)
var id pgtype.UUID
err := row.Scan(&id)
return id, err
}