package checksum import ( "context" "errors" "fmt" "log" "sync" "sync/atomic" "git.azzamo.net/Michilis/Spanglish/photo-api/internal/storage" "git.azzamo.net/Michilis/Spanglish/photo-api/internal/store" ) // Options configures a backfill run. type Options struct { // GalleryID limits the run to one gallery; empty means the whole library. GalleryID string // Concurrency is how many originals are hashed at a time. Concurrency int // DryRun reports what would be hashed without writing to the database. DryRun bool } // Result counts what a run did. Duplicates are rows whose bytes match an // earlier photo in the same gallery; they keep a NULL checksum and are listed // so an admin can decide what to do with them. type Result struct { Total int // rows without a checksum Hashed int // hashed and recorded (or, with DryRun, would be) Duplicates int Missing int // original object absent from storage Failed int } // Backfill hashes the stored original of every photo that has no checksum yet, // so duplicate detection also catches re-uploads of photos from before the // checksum column existed. It only ever writes the checksum column; no photo, // row or object is deleted. // // It is idempotent: rerunning it finds only what the previous run left behind. func Backfill(ctx context.Context, db *store.DB, st storage.Storage, opts Options) (Result, error) { if opts.Concurrency < 1 { opts.Concurrency = 4 } photos, err := db.PhotosMissingChecksum(ctx, opts.GalleryID) if err != nil { return Result{}, fmt.Errorf("list photos without a checksum: %w", err) } prefix := "" if opts.DryRun { prefix = "[dry-run] " } log.Printf("backfill-checksums: %s%d photos to hash, concurrency %d", prefix, len(photos), opts.Concurrency) var ( mu sync.Mutex res = Result{Total: len(photos)} done int64 jobs = make(chan store.Photo) wg sync.WaitGroup ) // Sequential hashing per worker, but the database update is what can // collide: two photos of the same gallery with identical bytes race, and // the unique index decides which one keeps the checksum. for i := 0; i < opts.Concurrency; i++ { wg.Add(1) go func() { defer wg.Done() for p := range jobs { if ctx.Err() != nil { return } sum, err := hashOriginal(ctx, st, p) if err == nil && !opts.DryRun { err = db.SetPhotoChecksum(ctx, p.ID, sum) } n := atomic.AddInt64(&done, 1) mu.Lock() switch { case errors.Is(err, storage.ErrNotExist): res.Missing++ log.Printf("backfill-checksums: [%d/%d] MISSING object, skipped: %s (photo %s)", n, len(photos), p.OriginalKey, p.ID) case store.IsUniqueViolation(err): res.Duplicates++ other := "an earlier photo" if existing, findErr := db.FindPhotoByChecksum(ctx, p.GalleryID, sum); findErr == nil { other = "photo " + existing.ID } log.Printf("backfill-checksums: [%d/%d] DUPLICATE: photo %s (%s) has the same content as %s "+ "in gallery %s — left without a checksum, delete it by hand if unwanted", n, len(photos), p.ID, p.OriginalFilename, other, p.GalleryID) case err != nil: res.Failed++ log.Printf("backfill-checksums: [%d/%d] FAILED %s (photo %s): %v", n, len(photos), p.OriginalKey, p.ID, err) default: res.Hashed++ } mu.Unlock() } }() } for _, p := range photos { select { case jobs <- p: case <-ctx.Done(): } if ctx.Err() != nil { break } } close(jobs) wg.Wait() log.Printf("backfill-checksums: %sdone — %d hashed, %d duplicates left unhashed, %d objects missing, %d failed", prefix, res.Hashed, res.Duplicates, res.Missing, res.Failed) if err := ctx.Err(); err != nil { return res, fmt.Errorf("interrupted after %d/%d photos: %w", res.Hashed, res.Total, err) } if res.Failed > 0 { return res, fmt.Errorf("%d of %d photos failed to hash (rerun to retry; already-hashed photos are skipped)", res.Failed, res.Total) } return res, nil } // hashOriginal reads a photo's stored original. Stat runs first because it is // the only call that reports a missing object as storage.ErrNotExist on both // backends — Open surfaces the driver's own error. func hashOriginal(ctx context.Context, st storage.Storage, p store.Photo) (string, error) { if _, err := st.Stat(ctx, p.OriginalKey); err != nil { return "", err } r, _, err := st.Open(ctx, p.OriginalKey) if err != nil { return "", err } defer r.Close() return Sum(r) }