Files
minstrel/internal/db/dbq/duplicates.sql.go
T
bvandeusenandClaude Opus 5.5 1e9408c835
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
feat(notifications): library health reaches admins, coalesced (#5341)
- 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>
2026-10-08 07:20:33 -04:00

476 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 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
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
}