MergeDuplicateGroup's loop body — repoint and copy every FK from a removed copy onto the survivor, inherit its MBID, delete its row, tidy an emptied album — moves into foldTrackInto, so the missing-pair pass can fold a stale missing row into its replacement with the same mechanics. The group lock, file removal and group bookkeeping stay in the merge. No behaviour change. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
345 lines
12 KiB
Go
345 lines
12 KiB
Go
package library
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log/slog"
|
|
"sort"
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgtype"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
|
|
"git.fabledsword.com/bvandeusen/minstrel/internal/coverart"
|
|
"git.fabledsword.com/bvandeusen/minstrel/internal/db/dbq"
|
|
syncpkg "git.fabledsword.com/bvandeusen/minstrel/internal/sync"
|
|
)
|
|
|
|
// Duplicate merge (M400 #3911).
|
|
|
|
// ErrDuplicateGroupNotPending means the group was already merged or dismissed,
|
|
// no longer exists, or no longer has two members to merge.
|
|
var ErrDuplicateGroupNotPending = errors.New("library: duplicate group is not pending")
|
|
|
|
// ErrSurvivorNotInGroup means the copy chosen to keep is not a member of the group.
|
|
var ErrSurvivorNotInGroup = errors.New("library: survivor is not a member of the group")
|
|
|
|
// MergedCopy is one copy a merge kept or removed.
|
|
type MergedCopy struct {
|
|
TrackID pgtype.UUID
|
|
FilePath string
|
|
TrackMbid *string
|
|
AlbumMbid *string
|
|
}
|
|
|
|
// MergeResult says what a merge did.
|
|
type MergeResult struct {
|
|
Tier string
|
|
Survivor MergedCopy
|
|
Removed []MergedCopy
|
|
|
|
// What moved onto the survivor — reported so the operator, and the audit
|
|
// log, can see that the history was kept rather than take it on trust.
|
|
PlayEvents int64
|
|
SkipEvents int64
|
|
Likes int // users whose like now sits on the survivor
|
|
PlaylistEntries int
|
|
|
|
DeletedAlbumIDs []pgtype.UUID
|
|
DeletedArtistIDs []pgtype.UUID
|
|
}
|
|
|
|
// MergeDuplicateGroup keeps one copy of a duplicate group and removes the rest,
|
|
// carrying everything the removed copies held onto the one kept.
|
|
//
|
|
// survivorID chooses the copy to keep; an invalid (zero) id takes the proposal
|
|
// from ProposeSurvivor.
|
|
//
|
|
// The danger this is built around: every table referencing tracks does so ON
|
|
// DELETE CASCADE, so deleting a duplicate's row outright silently destroys its
|
|
// likes, plays, playlist entries and tags. The merge moves all of that onto the
|
|
// survivor first, and only then deletes the now-empty row.
|
|
//
|
|
// It deletes the removed copies' FILES too, and first, before any row changes
|
|
// (#3918, note #3926). A merge that left the file behind would be undone by the
|
|
// next scan, which re-imports it as a new track with no history. If a file cannot
|
|
// be removed, the *FileRemoveError comes back and nothing in the database changes.
|
|
// With several copies to remove, one file may already be gone when a later one
|
|
// fails; that copy's row keeps all its history and is marked missing by the next
|
|
// scan, and retrying the merge picks up where it stopped.
|
|
//
|
|
// Everything else happens in one transaction, which holds a lock on the group so
|
|
// two merges of it cannot run at once. Sync changes for clients' caches are logged
|
|
// inside it, the way the playlists service logs its own.
|
|
func MergeDuplicateGroup(
|
|
ctx context.Context, pool *pgxpool.Pool, logger *slog.Logger, dataDir string,
|
|
groupID, survivorID pgtype.UUID,
|
|
) (MergeResult, error) {
|
|
if logger == nil {
|
|
logger = slog.Default()
|
|
}
|
|
tx, err := pool.Begin(ctx)
|
|
if err != nil {
|
|
return MergeResult{}, fmt.Errorf("begin merge: %w", err)
|
|
}
|
|
defer func() { _ = tx.Rollback(ctx) }()
|
|
tq := dbq.New(tx)
|
|
|
|
group, err := tq.LockDuplicateGroupForMerge(ctx, groupID)
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
return MergeResult{}, ErrDuplicateGroupNotPending
|
|
}
|
|
if err != nil {
|
|
return MergeResult{}, fmt.Errorf("lock duplicate group: %w", err)
|
|
}
|
|
if group.Status != "pending" {
|
|
return MergeResult{}, ErrDuplicateGroupNotPending
|
|
}
|
|
|
|
members, err := tq.ListDuplicateGroupMergeMembers(ctx, groupID)
|
|
if err != nil {
|
|
return MergeResult{}, fmt.Errorf("list group members: %w", err)
|
|
}
|
|
if len(members) < 2 {
|
|
return MergeResult{}, ErrDuplicateGroupNotPending
|
|
}
|
|
survivor, losers, err := splitSurvivor(members, survivorID)
|
|
if err != nil {
|
|
return MergeResult{}, err
|
|
}
|
|
|
|
for _, l := range losers {
|
|
if err := removeTrackFileOnDisk(l.FilePath); err != nil {
|
|
return MergeResult{}, err
|
|
}
|
|
}
|
|
|
|
res := MergeResult{Tier: group.Tier, Survivor: mergedCopyOf(survivor)}
|
|
likers := map[string]struct{}{}
|
|
changes := mergeChanges{}
|
|
|
|
for _, l := range losers {
|
|
f, err := foldTrackInto(ctx, tq, survivor.ID, l.ID, &changes)
|
|
if err != nil {
|
|
return MergeResult{}, err
|
|
}
|
|
res.PlayEvents += f.PlayEvents
|
|
res.SkipEvents += f.SkipEvents
|
|
res.PlaylistEntries += f.PlaylistEntries
|
|
for _, u := range f.Likers {
|
|
likers[u] = struct{}{}
|
|
}
|
|
if f.DeletedAlbumID != nil {
|
|
res.DeletedAlbumIDs = append(res.DeletedAlbumIDs, *f.DeletedAlbumID)
|
|
}
|
|
if f.DeletedArtistID != nil {
|
|
res.DeletedArtistIDs = append(res.DeletedArtistIDs, *f.DeletedArtistID)
|
|
}
|
|
res.Removed = append(res.Removed, mergedCopyOf(l))
|
|
}
|
|
res.Likes = len(likers)
|
|
|
|
marked, err := tq.MarkDuplicateGroupMerged(ctx, groupID)
|
|
if err != nil {
|
|
return MergeResult{}, fmt.Errorf("mark group merged: %w", err)
|
|
}
|
|
if marked != 1 {
|
|
return MergeResult{}, ErrDuplicateGroupNotPending
|
|
}
|
|
if err := changes.log(ctx, tx); err != nil {
|
|
return MergeResult{}, err
|
|
}
|
|
if err := tx.Commit(ctx); err != nil {
|
|
return MergeResult{}, fmt.Errorf("commit merge: %w", err)
|
|
}
|
|
|
|
// After commit, like DeleteTrackFile: a leftover art directory is only disk.
|
|
if dataDir != "" {
|
|
for _, artistID := range res.DeletedArtistIDs {
|
|
if err := coverart.CleanupArtistArt(dataDir, artistID); err != nil {
|
|
logger.Warn("duplicate merge: artist-art cleanup failed",
|
|
"artist_id", syncpkg.FormatUUID(artistID), "err", err)
|
|
}
|
|
}
|
|
}
|
|
return res, nil
|
|
}
|
|
|
|
// foldResult is what one fold moved onto the survivor.
|
|
type foldResult struct {
|
|
PlayEvents int64
|
|
SkipEvents int64
|
|
PlaylistEntries int
|
|
Likers []string // users whose like now sits on the survivor
|
|
|
|
DeletedAlbumID *pgtype.UUID
|
|
DeletedArtistID *pgtype.UUID
|
|
}
|
|
|
|
// foldTrackInto moves everything the loser row carries onto the survivor, then
|
|
// deletes the loser row and tidies an album or artist it leaves empty. It is the
|
|
// row half of a merge, shared by the duplicate merge and the missing-pair pass
|
|
// (M485): neither touches a file here, and both run it inside their own
|
|
// transaction. The sync changes it owes clients go into changes, which the
|
|
// caller logs once before committing.
|
|
func foldTrackInto(
|
|
ctx context.Context, tq *dbq.Queries, survivorID, loserID pgtype.UUID, changes *mergeChanges,
|
|
) (foldResult, error) {
|
|
var f foldResult
|
|
survivorKey := syncpkg.FormatUUID(survivorID)
|
|
loserKey := syncpkg.FormatUUID(loserID)
|
|
|
|
n, err := tq.MergeRepointPlayEvents(ctx, dbq.MergeRepointPlayEventsParams{SurvivorID: survivorID, LoserID: loserID})
|
|
if err != nil {
|
|
return f, fmt.Errorf("move play events: %w", err)
|
|
}
|
|
f.PlayEvents = n
|
|
n, err = tq.MergeRepointSkipEvents(ctx, dbq.MergeRepointSkipEventsParams{SurvivorID: survivorID, LoserID: loserID})
|
|
if err != nil {
|
|
return f, fmt.Errorf("move skip events: %w", err)
|
|
}
|
|
f.SkipEvents = n
|
|
if _, err := tq.MergeRepointContextualLikes(ctx, dbq.MergeRepointContextualLikesParams{SurvivorID: survivorID, LoserID: loserID}); err != nil {
|
|
return f, fmt.Errorf("move contextual likes: %w", err)
|
|
}
|
|
if _, err := tq.MergeRepointPlaybackErrors(ctx, dbq.MergeRepointPlaybackErrorsParams{SurvivorID: survivorID, LoserID: loserID}); err != nil {
|
|
return f, fmt.Errorf("move playback errors: %w", err)
|
|
}
|
|
if _, err := tq.MergeRepointLidarrRequests(ctx, dbq.MergeRepointLidarrRequestsParams{SurvivorID: survivorID, LoserID: loserID}); err != nil {
|
|
return f, fmt.Errorf("move lidarr request matches: %w", err)
|
|
}
|
|
|
|
playlists, err := tq.MergeRepointPlaylistTracks(ctx, dbq.MergeRepointPlaylistTracksParams{SurvivorID: survivorID, LoserID: loserID})
|
|
if err != nil {
|
|
return f, fmt.Errorf("move playlist entries: %w", err)
|
|
}
|
|
f.PlaylistEntries = len(playlists)
|
|
for _, pl := range playlists {
|
|
plKey := syncpkg.FormatUUID(pl)
|
|
changes.playlistDelete(syncpkg.EncodePlaylistTrackID(plKey, loserKey))
|
|
changes.playlistUpsert(syncpkg.EncodePlaylistTrackID(plKey, survivorKey))
|
|
}
|
|
|
|
users, err := tq.MergeCopyGeneralLikes(ctx, dbq.MergeCopyGeneralLikesParams{SurvivorID: survivorID, LoserID: loserID})
|
|
if err != nil {
|
|
return f, fmt.Errorf("move likes: %w", err)
|
|
}
|
|
for _, u := range users {
|
|
userKey := syncpkg.FormatUUID(u)
|
|
f.Likers = append(f.Likers, userKey)
|
|
changes.likeDelete(syncpkg.EncodeLikeID(userKey, loserKey))
|
|
changes.likeUpsert(syncpkg.EncodeLikeID(userKey, survivorKey))
|
|
}
|
|
|
|
if _, err := tq.MergeCopyTrackTags(ctx, dbq.MergeCopyTrackTagsParams{SurvivorID: survivorID, LoserID: loserID}); err != nil {
|
|
return f, fmt.Errorf("merge tags: %w", err)
|
|
}
|
|
if _, err := tq.MergeCopyTrackSimilarity(ctx, dbq.MergeCopyTrackSimilarityParams{SurvivorID: survivorID, LoserID: loserID}); err != nil {
|
|
return f, fmt.Errorf("merge similarity: %w", err)
|
|
}
|
|
if err := tq.MergeInheritTrackMbid(ctx, dbq.MergeInheritTrackMbidParams{SurvivorID: survivorID, LoserID: loserID}); err != nil {
|
|
return f, fmt.Errorf("inherit recording mbid: %w", err)
|
|
}
|
|
|
|
// Everything the loser carried now sits on the survivor, so the CASCADE
|
|
// this delete sets off has nothing left to destroy.
|
|
deleted, err := tq.DeleteTrack(ctx, loserID)
|
|
if err != nil {
|
|
return f, fmt.Errorf("delete merged copy: %w", err)
|
|
}
|
|
tidied, err := tidyEmptiedAlbum(ctx, tq, deleted.AlbumID)
|
|
if err != nil {
|
|
return f, err
|
|
}
|
|
f.DeletedAlbumID = tidied.AlbumID
|
|
f.DeletedArtistID = tidied.ArtistID
|
|
changes.trackDelete(loserKey)
|
|
return f, nil
|
|
}
|
|
|
|
// splitSurvivor separates the copy to keep from the copies to remove. An
|
|
// invalid survivorID takes ProposeSurvivor's choice.
|
|
func splitSurvivor(
|
|
members []dbq.ListDuplicateGroupMergeMembersRow, survivorID pgtype.UUID,
|
|
) (dbq.ListDuplicateGroupMergeMembersRow, []dbq.ListDuplicateGroupMergeMembersRow, error) {
|
|
want := ""
|
|
if survivorID.Valid {
|
|
want = syncpkg.FormatUUID(survivorID)
|
|
} else {
|
|
cands := make([]SurvivorCandidate, len(members))
|
|
for i, m := range members {
|
|
cands[i] = SurvivorCandidate{
|
|
TrackID: syncpkg.FormatUUID(m.ID), FileFormat: m.FileFormat, FileSize: m.FileSize, AddedAt: m.AddedAt.Time,
|
|
}
|
|
}
|
|
want, _ = ProposeSurvivor(cands)
|
|
}
|
|
|
|
var survivor dbq.ListDuplicateGroupMergeMembersRow
|
|
found := false
|
|
var losers []dbq.ListDuplicateGroupMergeMembersRow
|
|
for _, m := range members {
|
|
if syncpkg.FormatUUID(m.ID) == want {
|
|
survivor, found = m, true
|
|
continue
|
|
}
|
|
losers = append(losers, m)
|
|
}
|
|
if !found {
|
|
return dbq.ListDuplicateGroupMergeMembersRow{}, nil, ErrSurvivorNotInGroup
|
|
}
|
|
return survivor, losers, nil
|
|
}
|
|
|
|
func mergedCopyOf(m dbq.ListDuplicateGroupMergeMembersRow) MergedCopy {
|
|
return MergedCopy{TrackID: m.ID, FilePath: m.FilePath, TrackMbid: m.Mbid, AlbumMbid: m.AlbumMbid}
|
|
}
|
|
|
|
// mergeChanges collects the sync-log entries a merge owes clients' caches, each
|
|
// once: a user who liked two removed copies still gets one upsert for the
|
|
// survivor.
|
|
type mergeChanges struct {
|
|
tracks, likeDeletes, likeUpserts, playlistDeletes, playlistUpserts map[string]struct{}
|
|
}
|
|
|
|
func addTo(set *map[string]struct{}, id string) {
|
|
if *set == nil {
|
|
*set = map[string]struct{}{}
|
|
}
|
|
(*set)[id] = struct{}{}
|
|
}
|
|
|
|
func (c *mergeChanges) trackDelete(id string) { addTo(&c.tracks, id) }
|
|
func (c *mergeChanges) likeDelete(id string) { addTo(&c.likeDeletes, id) }
|
|
func (c *mergeChanges) likeUpsert(id string) { addTo(&c.likeUpserts, id) }
|
|
func (c *mergeChanges) playlistDelete(id string) { addTo(&c.playlistDeletes, id) }
|
|
func (c *mergeChanges) playlistUpsert(id string) { addTo(&c.playlistUpserts, id) }
|
|
|
|
func (c *mergeChanges) log(ctx context.Context, tx pgx.Tx) error {
|
|
for _, entry := range []struct {
|
|
kind syncpkg.EntityType
|
|
ids map[string]struct{}
|
|
op syncpkg.Op
|
|
}{
|
|
{syncpkg.EntityTrack, c.tracks, syncpkg.OpDelete},
|
|
{syncpkg.EntityLikeTrack, c.likeDeletes, syncpkg.OpDelete},
|
|
{syncpkg.EntityLikeTrack, c.likeUpserts, syncpkg.OpUpsert},
|
|
{syncpkg.EntityPlaylistTrack, c.playlistDeletes, syncpkg.OpDelete},
|
|
{syncpkg.EntityPlaylistTrack, c.playlistUpserts, syncpkg.OpUpsert},
|
|
} {
|
|
if len(entry.ids) == 0 {
|
|
continue
|
|
}
|
|
ids := make([]string, 0, len(entry.ids))
|
|
for id := range entry.ids {
|
|
ids = append(ids, id)
|
|
}
|
|
sort.Strings(ids)
|
|
if err := syncpkg.LogChanges(ctx, tx, entry.kind, ids, entry.op); err != nil {
|
|
return fmt.Errorf("log merge changes: %w", err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|