fix(downloads): enforce artifact links with an FK and make eviction leak-proof

Remediates review findings on the eviction fence:

- The EXISTS/NOT EXISTS fences read independent snapshots under READ
  COMMITTED and take no conflicting locks, so a concurrent evict/link
  pair can both commit (reproduced empirically in both interleavings),
  leaving a ready download pointing at a deleted artifact. New
  downloads_artifact_id_fkey (ON DELETE RESTRICT) is the enforcement
  layer: RI triggers take FOR KEY SHARE and re-check on a fresh snapshot.
  The fences stay as fast-path filters; 23503 maps to ErrArtifactEvicted
  on insert/replace and reads as 'not evictable' on delete. Terminal
  transitions null artifact_id so dead rows cannot pin artifacts; the
  migration heals legacy dangling/terminal links before the constraint.
- Row-first eviction could leak bytes when the unlink failed after the
  row was gone (no local orphan queue, no reaper; remote leaked when the
  node delete and orphan enqueue both failed). Cleanup intents are now
  recorded durably BEFORE the row delete — a new
  download_artifact_local_orphans queue mirrors the remote one — and a
  failed enqueue skips the candidate so the row survives for retry.
- Output paths now include the artifact id, so a replacement row can
  never claim an evicted row's path; the OutputPathInUse probe (an
  unindexed seq scan with a TOCTOU window) is gone.
- The failed-artifact hygiene sweep is fenced (DeleteIfEvictable):
  Ensure requeues failed rows, so the unfenced delete could remove an
  artifact a fresh download had just linked, stranding it in preparing.
- ReplaceManagedEntry drops the non-atomic artifact probe; zero rows now
  unambiguously means revision/identity conflict and self-heals via
  GetManagedEntry. The replace retry refreshes its row snapshot so the
  second attempt's revision fence can succeed.
- Reuse/cleanups: shared withEnsuredArtifact retry helper, downloadColumns
  interpolated into all insert sites, execInsertDownload shared by
  Create/CreateBatch, EnsureDevice hoisted out of the retry loop,
  HasActiveLink and unconditional DeleteArtifact deleted, distinct log
  messages for the two sweep failure paths.

DB-backed tests cover the FK both directions, the requeued-artifact sweep
fence, local orphan round-trip, terminal unpinning, revision-conflict
self-healing, and a 100-round concurrent evict-vs-link race asserting no
dangling links.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Quick104
2026-08-13 12:32:11 -04:00
co-authored by Claude Fable 5
parent 1ae35c244f
commit 80397a5eb3
11 changed files with 946 additions and 207 deletions
+8 -4
View File
@@ -74,12 +74,16 @@ func effectiveArtifactDir(artifactDir, transcodeDir string) string {
return config.EffectiveDownloadArtifactDir(artifactDir, transcodeDir)
}
// artifactOutputPath derives a deterministic output path from
// (media_file_id, format, params_hash) so a reclaimed job targets the same file.
func artifactOutputPath(dir string, mediaFileID int, format, hash string) string {
// artifactOutputPath derives an output path from
// (media_file_id, format, params_hash, artifact_id). It is deterministic for a
// given artifact row — a reclaimed job targets the same file — while the
// artifact id keeps a REPLACEMENT row's path distinct from its predecessor's.
// Without that, an evicted row's bytes could be unlinked after a successor with
// the same dedup key had already claimed the identical path.
func artifactOutputPath(dir string, mediaFileID int, format, hash, artifactID string) string {
short := hash
if len(short) > 16 {
short = short[:16]
}
return filepath.Join(dir, fmt.Sprintf("%d_%s_%s.mp4", mediaFileID, format, short))
return filepath.Join(dir, fmt.Sprintf("%d_%s_%s_%s.mp4", mediaFileID, format, short, artifactID))
}
+122 -44
View File
@@ -408,21 +408,6 @@ func (r *ArtifactRepository) TotalReadyBytes(ctx context.Context) (int64, error)
return total, nil
}
// HasActiveLink reports whether any valid download row — managed or ephemeral
// (device-less web) — still references the artifact. Completed rows are
// retained because they remain re-downloadable handles; evicting an artifact a
// live row references would 404 a download the API advertises as servable.
func (r *ArtifactRepository) HasActiveLink(ctx context.Context, artifactID string) (bool, error) {
var exists bool
if err := r.pool.QueryRow(ctx,
`SELECT `+activeArtifactLinkSQL,
artifactID,
).Scan(&exists); err != nil {
return false, fmt.Errorf("checking artifact links: %w", err)
}
return exists, nil
}
// ListFailedBefore returns terminally-failed artifacts cold since cutoff
// (last_used_at). Their linked downloads were already flipped to 'failed' by
// reconciliation, so the rows serve nothing and only block re-attempts.
@@ -452,46 +437,124 @@ func (r *ArtifactRepository) ListUnlinkedReadyBefore(ctx context.Context, cutoff
return scanArtifacts(rows)
}
// DeleteArtifact removes an artifact row.
func (r *ArtifactRepository) DeleteArtifact(ctx context.Context, id string) error {
_, err := r.pool.Exec(ctx, `DELETE FROM download_artifacts WHERE id = $1`, id)
if err != nil {
return fmt.Errorf("deleting artifact: %w", err)
}
return nil
}
// DeleteReadyIfEvictable atomically deletes a ready artifact only when no
// active download row references it. Check-then-delete eviction races a new
// download linking the same id; this fence makes that outcome impossible.
// download linking the same id; the NOT-EXISTS predicate is the fast path and
// downloads_artifact_id_fkey (ON DELETE RESTRICT) is the enforcement backstop —
// a concurrent link that commits after this statement's snapshot makes the RI
// trigger raise 23503, which reads as "not evictable".
func (r *ArtifactRepository) DeleteReadyIfEvictable(ctx context.Context, id string) (bool, error) {
tag, err := r.pool.Exec(ctx,
`DELETE FROM download_artifacts
WHERE id = $1 AND status = 'ready'
AND NOT `+activeArtifactLinkSQL,
id,
)
return r.deleteIfEvictable(ctx, id, true)
}
// DeleteIfEvictable is DeleteReadyIfEvictable without the ready-status
// predicate: the failed-artifact sweep uses it so a row Ensure requeued (and a
// download then linked) is never removed out from under that link.
func (r *ArtifactRepository) DeleteIfEvictable(ctx context.Context, id string) (bool, error) {
return r.deleteIfEvictable(ctx, id, false)
}
func (r *ArtifactRepository) deleteIfEvictable(ctx context.Context, id string, requireReady bool) (bool, error) {
query := `DELETE FROM download_artifacts WHERE id = $1`
if requireReady {
query += ` AND status = 'ready'`
}
query += ` AND NOT ` + activeArtifactLinkSQL
tag, err := r.pool.Exec(ctx, query, id)
if err != nil {
// A download row linked the artifact concurrently: not evictable.
if isForeignKeyViolation(err, downloadsArtifactFKConstraint) {
return false, nil
}
return false, fmt.Errorf("evicting unlinked artifact: %w", err)
}
return tag.RowsAffected() == 1, nil
}
// OutputPathInUse reports whether any remaining artifact row still claims path.
// Eviction deletes the DB row before bytes; a concurrent EnsureQueued may have
// already reused the deterministic local path for a replacement job.
func (r *ArtifactRepository) OutputPathInUse(ctx context.Context, path string) (bool, error) {
if path == "" {
return false, nil
// LocalArtifactOrphan is a filesystem cleanup candidate recorded before its
// artifact row is deleted, so a crash or unlink failure between the row delete
// and the byte delete can never leak the file.
type LocalArtifactOrphan struct {
ID int64
DownloadArtifactID string
OutputPath string
Attempts int
}
// EnqueueLocalOrphan durably records a local artifact path for deletion. Paths
// are unique per artifact id, so a conflicting row already covers the same file.
func (r *ArtifactRepository) EnqueueLocalOrphan(ctx context.Context, artifactID, outputPath string) error {
if artifactID == "" || outputPath == "" {
return nil
}
var exists bool
if err := r.pool.QueryRow(ctx,
`SELECT EXISTS(SELECT 1 FROM download_artifacts WHERE output_path = $1)`,
path,
).Scan(&exists); err != nil {
return false, fmt.Errorf("checking artifact output path: %w", err)
_, err := r.pool.Exec(ctx,
`INSERT INTO download_artifact_local_orphans (download_artifact_id, output_path)
VALUES ($1, $2)
ON CONFLICT (output_path) DO NOTHING`,
artifactID, outputPath,
)
if err != nil {
return fmt.Errorf("enqueueing local artifact cleanup: %w", err)
}
return exists, nil
return nil
}
// ListLocalOrphansDue returns due local cleanup candidates, oldest first.
func (r *ArtifactRepository) ListLocalOrphansDue(ctx context.Context, limit int) ([]LocalArtifactOrphan, error) {
if limit <= 0 {
limit = 100
}
rows, err := r.pool.Query(ctx,
`SELECT id, download_artifact_id, output_path, attempts
FROM download_artifact_local_orphans
WHERE next_retry_at IS NULL OR next_retry_at <= now()
ORDER BY attempts, created_at, id
LIMIT $1`, limit)
if err != nil {
return nil, fmt.Errorf("listing local artifact cleanup queue: %w", err)
}
defer rows.Close()
out := make([]LocalArtifactOrphan, 0, limit)
for rows.Next() {
var orphan LocalArtifactOrphan
if err := rows.Scan(&orphan.ID, &orphan.DownloadArtifactID, &orphan.OutputPath, &orphan.Attempts); err != nil {
return nil, fmt.Errorf("scanning local artifact cleanup row: %w", err)
}
out = append(out, orphan)
}
return out, rows.Err()
}
// DeleteLocalOrphan clears a local cleanup row by id.
func (r *ArtifactRepository) DeleteLocalOrphan(ctx context.Context, id int64) error {
if _, err := r.pool.Exec(ctx, `DELETE FROM download_artifact_local_orphans WHERE id = $1`, id); err != nil {
return fmt.Errorf("deleting local artifact cleanup row: %w", err)
}
return nil
}
// DeleteLocalOrphanByPath clears a local cleanup row after the inline unlink
// that followed its enqueue succeeded.
func (r *ArtifactRepository) DeleteLocalOrphanByPath(ctx context.Context, outputPath string) error {
if outputPath == "" {
return nil
}
if _, err := r.pool.Exec(ctx, `DELETE FROM download_artifact_local_orphans WHERE output_path = $1`, outputPath); err != nil {
return fmt.Errorf("deleting local artifact cleanup row by path: %w", err)
}
return nil
}
// BumpLocalOrphanRetry defers a failed unlink behind the same 2-minute backoff
// the remote queue uses.
func (r *ArtifactRepository) BumpLocalOrphanRetry(ctx context.Context, id int64) error {
if _, err := r.pool.Exec(ctx,
`UPDATE download_artifact_local_orphans
SET attempts = attempts + 1, next_retry_at = now() + interval '2 minutes'
WHERE id = $1`, id); err != nil {
return fmt.Errorf("deferring local artifact cleanup row: %w", err)
}
return nil
}
// EnqueueRemoteOrphan durably records an abandoned attempt before its owning
@@ -595,6 +658,21 @@ func (r *ArtifactRepository) PrepareRemoteOrphanCleanup(ctx context.Context, orp
return owned, true, nil
}
// DeleteRemoteOrphanByLocator clears a cleanup intent whose remote delete
// already succeeded inline.
func (r *ArtifactRepository) DeleteRemoteOrphanByLocator(ctx context.Context, nodeID int, artifactID string) error {
if nodeID <= 0 || artifactID == "" {
return nil
}
if _, err := r.pool.Exec(ctx,
`DELETE FROM download_artifact_orphans WHERE origin_node_id = $1 AND origin_artifact_id = $2`,
nodeID, artifactID,
); err != nil {
return fmt.Errorf("deleting remote artifact cleanup row by locator: %w", err)
}
return nil
}
func (r *ArtifactRepository) DeleteRemoteOrphan(ctx context.Context, id int64) error {
if _, err := r.pool.Exec(ctx, `DELETE FROM download_artifact_orphans WHERE id = $1`, id); err != nil {
return fmt.Errorf("deleting remote artifact cleanup row: %w", err)
+367 -10
View File
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"fmt"
"sync"
"testing"
"time"
@@ -48,6 +49,9 @@ func newArtifactTestRepo(t *testing.T) (*ArtifactRepository, *pgxpool.Pool, int)
t.Fatalf("seed media file: %v", err)
}
t.Cleanup(func() {
// downloads_artifact_id_fkey is ON DELETE RESTRICT: unpin any surviving
// links before removing the artifact rows they reference.
_, _ = pool.Exec(ctx, `UPDATE downloads SET artifact_id = NULL WHERE media_file_id = $1`, fileID)
_, _ = pool.Exec(ctx, `DELETE FROM download_artifacts WHERE media_file_id = $1`, fileID)
_, _ = pool.Exec(ctx, `DELETE FROM media_files WHERE id = $1`, fileID)
_, _ = pool.Exec(ctx, `DELETE FROM media_folders WHERE id = $1`, folderID)
@@ -718,10 +722,10 @@ func TestListRemoteOrphansDueIsFairAcrossOrigins(t *testing.T) {
}
}
// TestHasActiveLinkCoversEphemeralRows pins the eviction guard: an ephemeral
// TestEvictionFenceCoversEphemeralRows pins the eviction guard: an ephemeral
// (device-less web) download row must protect its artifact from LRU cleanup
// exactly like a managed row does, and terminal rows must not.
func TestHasActiveLinkCoversEphemeralRows(t *testing.T) {
func TestEvictionFenceCoversEphemeralRows(t *testing.T) {
repo, pool, fileID := newArtifactTestRepo(t)
ctx := context.Background()
@@ -729,6 +733,7 @@ func TestHasActiveLinkCoversEphemeralRows(t *testing.T) {
if _, _, err := repo.EnsureQueued(ctx, art); err != nil {
t.Fatalf("ensure artifact: %v", err)
}
markArtifactReady(t, pool, art.ID)
var userID int
if err := pool.QueryRow(ctx,
@@ -754,26 +759,146 @@ func TestHasActiveLinkCoversEphemeralRows(t *testing.T) {
t.Fatalf("create ephemeral download: %v", err)
}
active, err := repo.HasActiveLink(ctx, art.ID)
deleted, err := repo.DeleteReadyIfEvictable(ctx, art.ID)
if err != nil {
t.Fatalf("HasActiveLink: %v", err)
t.Fatalf("DeleteReadyIfEvictable with ephemeral link: %v", err)
}
if !active {
if deleted {
t.Fatal("ephemeral ready row must protect its artifact from eviction")
}
if _, err := pool.Exec(ctx, `UPDATE downloads SET status = 'cancelled' WHERE id = $1`, dlID); err != nil {
if _, err := pool.Exec(ctx, `UPDATE downloads SET status = 'cancelled', artifact_id = NULL WHERE id = $1`, dlID); err != nil {
t.Fatalf("cancel download: %v", err)
}
active, err = repo.HasActiveLink(ctx, art.ID)
deleted, err = repo.DeleteReadyIfEvictable(ctx, art.ID)
if err != nil {
t.Fatalf("HasActiveLink after cancel: %v", err)
t.Fatalf("DeleteReadyIfEvictable after cancel: %v", err)
}
if active {
if !deleted {
t.Fatal("terminal-only links must not protect an artifact")
}
}
// TestDeleteIfEvictableFencesRequeuedArtifact covers the failed-artifact sweep:
// a row Ensure requeued (so it is no longer 'ready') and a download then linked
// must survive the sweep, and become removable once the link goes terminal.
func TestDeleteIfEvictableFencesRequeuedArtifact(t *testing.T) {
repo, pool, fileID := newArtifactTestRepo(t)
ctx := context.Background()
art := newArtifact(t, fileID, fmt.Sprintf("hash-requeued-%d", time.Now().UnixNano()))
if _, _, err := repo.EnsureQueued(ctx, art); err != nil {
t.Fatalf("ensure artifact: %v", err)
}
userID := seedDownloadUser(t, pool)
dlRepo := NewRepository(pool)
now := time.Now()
dlID := fmt.Sprintf("dl-requeued-%d", now.UnixNano())
if err := dlRepo.Create(ctx, &Download{
ID: dlID, UserID: userID, MediaFileID: fileID,
ContentID: fmt.Sprintf("requeued-content-%d", now.UnixNano()),
Kind: KindQueued, Status: StatusPreparing, Format: FormatTranscode,
ArtifactID: art.ID, FileSize: 1024, CreatedAt: now, UpdatedAt: now,
}); err != nil {
t.Fatalf("create preparing download: %v", err)
}
deleted, err := repo.DeleteIfEvictable(ctx, art.ID)
if err != nil {
t.Fatalf("DeleteIfEvictable linked: %v", err)
}
if deleted {
t.Fatal("a linked non-ready artifact must not be swept")
}
if _, err := repo.GetByID(ctx, art.ID); err != nil {
t.Fatalf("linked artifact missing after sweep attempt: %v", err)
}
if _, err := pool.Exec(ctx, `UPDATE downloads SET status = 'failed', artifact_id = NULL WHERE id = $1`, dlID); err != nil {
t.Fatalf("fail download: %v", err)
}
deleted, err = repo.DeleteIfEvictable(ctx, art.ID)
if err != nil {
t.Fatalf("DeleteIfEvictable after fail: %v", err)
}
if !deleted {
t.Fatal("an unlinked non-ready artifact must be swept")
}
}
// TestLocalOrphanQueueRoundTrip covers the durable local cleanup intent: an
// enqueue is idempotent per path, becomes due immediately, defers on a bumped
// retry, and clears by id and by path.
func TestLocalOrphanQueueRoundTrip(t *testing.T) {
repo, pool, fileID := newArtifactTestRepo(t)
ctx := context.Background()
var present *string
if err := pool.QueryRow(ctx, `SELECT to_regclass('public.download_artifact_local_orphans')::text`).Scan(&present); err != nil {
t.Fatalf("check download_artifact_local_orphans: %v", err)
}
if present == nil {
t.Skip("download_artifact_local_orphans migration has not been applied")
}
art := newArtifact(t, fileID, fmt.Sprintf("hash-local-orphan-%d", time.Now().UnixNano()))
path := art.OutputPath
t.Cleanup(func() {
_, _ = pool.Exec(context.Background(), `DELETE FROM download_artifact_local_orphans WHERE output_path = $1`, path)
})
if err := repo.EnqueueLocalOrphan(ctx, art.ID, path); err != nil {
t.Fatalf("EnqueueLocalOrphan: %v", err)
}
if err := repo.EnqueueLocalOrphan(ctx, art.ID, path); err != nil {
t.Fatalf("EnqueueLocalOrphan (repeat): %v", err)
}
due, err := repo.ListLocalOrphansDue(ctx, 500)
if err != nil {
t.Fatalf("ListLocalOrphansDue: %v", err)
}
mine := localOrphansForPath(due, path)
if len(mine) != 1 {
t.Fatalf("due local orphans for path = %d, want exactly 1 (enqueue must dedupe)", len(mine))
}
if mine[0].DownloadArtifactID != art.ID {
t.Fatalf("orphan artifact id = %q, want %q", mine[0].DownloadArtifactID, art.ID)
}
if err := repo.BumpLocalOrphanRetry(ctx, mine[0].ID); err != nil {
t.Fatalf("BumpLocalOrphanRetry: %v", err)
}
due, err = repo.ListLocalOrphansDue(ctx, 500)
if err != nil {
t.Fatalf("ListLocalOrphansDue after bump: %v", err)
}
if len(localOrphansForPath(due, path)) != 0 {
t.Fatal("a bumped local orphan must not be due again immediately")
}
if err := repo.DeleteLocalOrphanByPath(ctx, path); err != nil {
t.Fatalf("DeleteLocalOrphanByPath: %v", err)
}
var remaining int
if err := pool.QueryRow(ctx,
`SELECT count(*) FROM download_artifact_local_orphans WHERE output_path = $1`, path,
).Scan(&remaining); err != nil {
t.Fatalf("count local orphans: %v", err)
}
if remaining != 0 {
t.Fatalf("local orphan rows after delete = %d, want 0", remaining)
}
}
func localOrphansForPath(orphans []LocalArtifactOrphan, path string) []LocalArtifactOrphan {
filtered := make([]LocalArtifactOrphan, 0, 1)
for _, orphan := range orphans {
if orphan.OutputPath == path {
filtered = append(filtered, orphan)
}
}
return filtered
}
func markArtifactReady(t *testing.T, pool *pgxpool.Pool, id string) {
t.Helper()
if _, err := pool.Exec(context.Background(),
@@ -856,7 +981,8 @@ func TestDeleteReadyIfEvictableFencesConcurrentLink(t *testing.T) {
t.Fatalf("linked artifact missing after eviction attempt: %v", err)
}
if _, err := pool.Exec(ctx, `UPDATE downloads SET status = 'cancelled' WHERE artifact_id = $1`, art2.ID); err != nil {
// Terminal transitions drop the link, as CancelByID and the fail paths do.
if _, err := pool.Exec(ctx, `UPDATE downloads SET status = 'cancelled', artifact_id = NULL WHERE artifact_id = $1`, art2.ID); err != nil {
t.Fatalf("cancel download: %v", err)
}
deleted, err = repo.DeleteReadyIfEvictable(ctx, art2.ID)
@@ -882,3 +1008,234 @@ func TestCreateOriginalDownloadDoesNotRequireArtifact(t *testing.T) {
t.Fatalf("Create original download: %v", err)
}
}
// TestCreateAgainstMissingArtifactIsEvicted pins the enforcement layer itself:
// a download that names an artifact id with no row must be rejected as evicted,
// whichever mechanism catches it (the EXISTS fast path or the foreign key).
func TestCreateAgainstMissingArtifactIsEvicted(t *testing.T) {
_, pool, fileID := newArtifactTestRepo(t)
ctx := context.Background()
userID := seedDownloadUser(t, pool)
now := time.Now()
missingID, err := idgen.NextID()
if err != nil {
t.Fatalf("id: %v", err)
}
d := &Download{
ID: fmt.Sprintf("dl-missing-%d", now.UnixNano()), UserID: userID, MediaFileID: fileID,
ContentID: fmt.Sprintf("missing-content-%d", now.UnixNano()),
Kind: KindQueued, Status: StatusReady, Format: FormatTranscode,
ArtifactID: missingID, FileSize: 1024, CreatedAt: now, UpdatedAt: now,
}
if err := NewRepository(pool).Create(ctx, d); !errors.Is(err, ErrArtifactEvicted) {
t.Fatalf("Create against missing artifact = %v, want ErrArtifactEvicted", err)
}
if err := NewRepository(pool).CreateBatch(ctx, []*Download{d}); !errors.Is(err, ErrArtifactEvicted) {
t.Fatalf("CreateBatch against missing artifact = %v, want ErrArtifactEvicted", err)
}
}
// TestTerminalTransitionsUnpinArtifact pins the ON DELETE RESTRICT contract:
// cancelling or failing a download must drop its artifact link, otherwise the
// artifact stays pinned against eviction forever.
func TestTerminalTransitionsUnpinArtifact(t *testing.T) {
repo, pool, fileID := newArtifactTestRepo(t)
ctx := context.Background()
dlRepo := NewRepository(pool)
userID := seedDownloadUser(t, pool)
now := time.Now()
cancelArt := newArtifact(t, fileID, fmt.Sprintf("hash-cancel-%d", now.UnixNano()))
if _, _, err := repo.EnsureQueued(ctx, cancelArt); err != nil {
t.Fatalf("ensure cancel artifact: %v", err)
}
markArtifactReady(t, pool, cancelArt.ID)
cancelID := fmt.Sprintf("dl-cancel-%d", now.UnixNano())
if err := dlRepo.Create(ctx, &Download{
ID: cancelID, UserID: userID, MediaFileID: fileID,
ContentID: fmt.Sprintf("cancel-content-%d", now.UnixNano()),
Kind: KindQueued, Status: StatusQueued, Format: FormatTranscode,
ArtifactID: cancelArt.ID, FileSize: 1024, CreatedAt: now, UpdatedAt: now,
}); err != nil {
t.Fatalf("create cancellable download: %v", err)
}
if err := dlRepo.CancelByID(ctx, cancelID, userID); err != nil {
t.Fatalf("CancelByID: %v", err)
}
cancelled, err := dlRepo.GetByID(ctx, cancelID)
if err != nil {
t.Fatalf("get cancelled download: %v", err)
}
if cancelled.ArtifactID != "" {
t.Fatalf("cancelled download artifact_id = %q, want empty", cancelled.ArtifactID)
}
deleted, err := repo.DeleteReadyIfEvictable(ctx, cancelArt.ID)
if err != nil {
t.Fatalf("DeleteReadyIfEvictable after cancel: %v", err)
}
if !deleted {
t.Fatal("a cancelled download must not pin its artifact")
}
failArt := newArtifact(t, fileID, fmt.Sprintf("hash-fail-%d", now.UnixNano()))
if _, _, err := repo.EnsureQueued(ctx, failArt); err != nil {
t.Fatalf("ensure fail artifact: %v", err)
}
failID := fmt.Sprintf("dl-fail-%d", now.UnixNano())
if err := dlRepo.Create(ctx, &Download{
ID: failID, UserID: userID, MediaFileID: fileID,
ContentID: fmt.Sprintf("fail-content-%d", now.UnixNano()),
Kind: KindQueued, Status: StatusPreparing, Format: FormatTranscode,
ArtifactID: failArt.ID, FileSize: 1024, CreatedAt: now, UpdatedAt: now,
}); err != nil {
t.Fatalf("create preparing download: %v", err)
}
flipped, err := dlRepo.MarkLinkedDownloadsFailed(ctx, failArt.ID, "encode failed")
if err != nil {
t.Fatalf("MarkLinkedDownloadsFailed: %v", err)
}
if len(flipped) != 1 {
t.Fatalf("flipped rows = %d, want 1", len(flipped))
}
if flipped[0].ArtifactID != "" {
t.Fatalf("failed download artifact_id = %q, want empty", flipped[0].ArtifactID)
}
deleted, err = repo.DeleteIfEvictable(ctx, failArt.ID)
if err != nil {
t.Fatalf("DeleteIfEvictable after fail: %v", err)
}
if !deleted {
t.Fatal("a failed download must not pin its artifact")
}
}
// TestReconcileFailedDownloadsUnpinArtifact covers the reconcile-to-failed
// update: it must both preserve the artifact's error message (read from the OLD
// row) and clear the link in the same statement.
func TestReconcileFailedDownloadsUnpinArtifact(t *testing.T) {
repo, pool, fileID := newArtifactTestRepo(t)
ctx := context.Background()
dlRepo := NewRepository(pool)
userID := seedDownloadUser(t, pool)
now := time.Now()
art := newArtifact(t, fileID, fmt.Sprintf("hash-reconcile-%d", now.UnixNano()))
if _, _, err := repo.EnsureQueued(ctx, art); err != nil {
t.Fatalf("ensure artifact: %v", err)
}
if _, err := pool.Exec(ctx,
`UPDATE download_artifacts SET status = 'failed', error_message = 'ffmpeg exploded' WHERE id = $1`, art.ID,
); err != nil {
t.Fatalf("mark artifact failed: %v", err)
}
dlID := fmt.Sprintf("dl-reconcile-%d", now.UnixNano())
if err := dlRepo.Create(ctx, &Download{
ID: dlID, UserID: userID, MediaFileID: fileID,
ContentID: fmt.Sprintf("reconcile-content-%d", now.UnixNano()),
Kind: KindQueued, Status: StatusPreparing, Format: FormatTranscode,
ArtifactID: art.ID, FileSize: 1024, CreatedAt: now, UpdatedAt: now,
}); err != nil {
t.Fatalf("create preparing download: %v", err)
}
if _, _, err := dlRepo.ReconcileLinkedDownloads(ctx); err != nil {
t.Fatalf("ReconcileLinkedDownloads: %v", err)
}
row, err := dlRepo.GetByID(ctx, dlID)
if err != nil {
t.Fatalf("get reconciled download: %v", err)
}
if row.Status != StatusFailed {
t.Fatalf("reconciled status = %q, want failed", row.Status)
}
if row.ErrorMessage != "ffmpeg exploded" {
t.Fatalf("reconciled error_message = %q, want the artifact's message", row.ErrorMessage)
}
if row.ArtifactID != "" {
t.Fatalf("reconciled artifact_id = %q, want empty", row.ArtifactID)
}
}
// TestEvictionVersusLinkRace is the invariant test the fence alone could not
// satisfy: many rounds of a genuine two-goroutine race between linking an
// artifact and evicting it, asserting after every round that no download row
// points at a missing artifact and that exactly one side lost.
func TestEvictionVersusLinkRace(t *testing.T) {
repo, pool, fileID := newArtifactTestRepo(t)
ctx := context.Background()
dlRepo := NewRepository(pool)
userID := seedDownloadUser(t, pool)
const rounds = 100
for i := 0; i < rounds; i++ {
art := newArtifact(t, fileID, fmt.Sprintf("hash-race-%d-%d", time.Now().UnixNano(), i))
if _, _, err := repo.EnsureQueued(ctx, art); err != nil {
t.Fatalf("round %d: ensure artifact: %v", i, err)
}
markArtifactReady(t, pool, art.ID)
now := time.Now()
d := &Download{
ID: fmt.Sprintf("dl-race-%d-%d", now.UnixNano(), i), UserID: userID, MediaFileID: fileID,
ContentID: fmt.Sprintf("race-content-%d-%d", now.UnixNano(), i),
Kind: KindQueued, Status: StatusReady, Format: FormatTranscode,
ArtifactID: art.ID, FileSize: 1024, CreatedAt: now, UpdatedAt: now,
}
start := make(chan struct{})
var wg sync.WaitGroup
var createErr, deleteErr error
var evicted bool
wg.Add(2)
go func() {
defer wg.Done()
<-start
createErr = dlRepo.Create(ctx, d)
}()
go func() {
defer wg.Done()
<-start
evicted, deleteErr = repo.DeleteReadyIfEvictable(ctx, art.ID)
}()
close(start)
wg.Wait()
if deleteErr != nil {
t.Fatalf("round %d: DeleteReadyIfEvictable: %v", i, deleteErr)
}
if createErr != nil && !errors.Is(createErr, ErrArtifactEvicted) {
t.Fatalf("round %d: Create: %v", i, createErr)
}
// The invariant: never a download row whose artifact row is gone.
var dangling int
if err := pool.QueryRow(ctx,
`SELECT count(*) FROM downloads d
WHERE d.artifact_id IS NOT NULL
AND NOT EXISTS (SELECT 1 FROM download_artifacts a WHERE a.id = d.artifact_id)`,
).Scan(&dangling); err != nil {
t.Fatalf("round %d: dangling check: %v", i, err)
}
if dangling != 0 {
t.Fatalf("round %d: %d download rows point at a deleted artifact", i, dangling)
}
// And the eviction direction: exactly one side may win.
linked := createErr == nil
if linked && evicted {
t.Fatalf("round %d: both the link and the eviction succeeded", i)
}
if linked {
if _, err := repo.GetByID(ctx, art.ID); err != nil {
t.Fatalf("round %d: linked artifact missing: %v", i, err)
}
}
if _, err := pool.Exec(ctx, `DELETE FROM downloads WHERE id = $1`, d.ID); err != nil {
t.Fatalf("round %d: cleanup download: %v", i, err)
}
if _, err := pool.Exec(ctx, `DELETE FROM download_artifacts WHERE id = $1`, art.ID); err != nil {
t.Fatalf("round %d: cleanup artifact: %v", i, err)
}
}
}
+11 -2
View File
@@ -39,8 +39,8 @@ func TestParamsHashStableAndDistinct(t *testing.T) {
}
func TestArtifactOutputPathDeterministic(t *testing.T) {
p1 := artifactOutputPath("/var/artifacts", 42, "transcode", "abcdef0123456789deadbeef")
p2 := artifactOutputPath("/var/artifacts", 42, "transcode", "abcdef0123456789deadbeef")
p1 := artifactOutputPath("/var/artifacts", 42, "transcode", "abcdef0123456789deadbeef", "art-1")
p2 := artifactOutputPath("/var/artifacts", 42, "transcode", "abcdef0123456789deadbeef", "art-1")
if p1 != p2 {
t.Fatalf("output path not deterministic: %q != %q", p1, p2)
}
@@ -50,6 +50,15 @@ func TestArtifactOutputPathDeterministic(t *testing.T) {
if !strings.Contains(p1, "42_transcode_") {
t.Fatalf("output path missing identity components: %q", p1)
}
// A replacement row for the same dedup key must NOT reuse the evicted row's
// path, or its bytes could be unlinked by the predecessor's cleanup.
p3 := artifactOutputPath("/var/artifacts", 42, "transcode", "abcdef0123456789deadbeef", "art-2")
if p3 == p1 {
t.Fatalf("distinct artifact ids share an output path: %q", p3)
}
if !strings.Contains(p3, "art-2") {
t.Fatalf("output path missing artifact id: %q", p3)
}
}
type lifecycleTestPreparer struct {
+152 -62
View File
@@ -243,7 +243,7 @@ func (m *ArtifactManager) Ensure(ctx context.Context, file *models.MediaFile, fo
Resolution: target.Resolution,
AudioTrackIndex: target.AudioTrackIndex,
TargetBitrateKbps: target.TargetBitrateKbps,
OutputPath: artifactOutputPath(m.artifactDir(), file.ID, format, hash),
OutputPath: artifactOutputPath(m.artifactDir(), file.ID, format, hash, id),
MaxAttempts: artifactMaxAttempts,
}
row, created, err := m.repo.EnsureQueued(ctx, a)
@@ -793,6 +793,7 @@ const (
// all terminal are evictable.
func (m *ArtifactManager) Cleanup(ctx context.Context) error {
m.cleanupRemoteOrphans(ctx)
m.cleanupLocalOrphans(ctx)
m.sweepStale(ctx)
budget := m.downloadConfig().ArtifactMaxBytes
if budget <= 0 {
@@ -813,12 +814,18 @@ func (m *ArtifactManager) Cleanup(ctx context.Context) error {
if total <= budget {
break
}
// Delete the row first, gated on the same active-link predicate
// HasActiveLink uses. Checking then deleting bytes then deleting the
// row lets a concurrent Ensure+Create land a ready download on an id
// that this pass then removes — the download stays advertised and
// never heals. The atomic delete closes that window; Create refuses
// to link a row that no longer exists.
// Record the byte-cleanup intent BEFORE deleting the row: once the row
// is gone nothing else remembers the file, so a crash or unlink failure
// after the delete would leak it. A failed enqueue skips the candidate
// entirely — the row survives and the next pass retries.
if !m.enqueueEvictionCleanup(ctx, a) {
continue
}
// Delete the row first, gated on the active-link predicate. Checking,
// then deleting bytes, then deleting the row lets a concurrent
// Ensure+Create land a ready download on an id that this pass then
// removes — the download stays advertised and never heals. The fenced
// delete plus downloads_artifact_id_fkey closes that window.
deleted, err := m.repo.DeleteReadyIfEvictable(ctx, a.ID)
if err != nil {
slog.WarnContext(ctx, "artifact eviction fence failed", "component", "downloads", "artifact_id", a.ID, "error", err)
@@ -827,15 +834,100 @@ func (m *ArtifactManager) Cleanup(ctx context.Context) error {
if !deleted {
continue
}
if !m.deleteEvictedArtifactBytes(ctx, a) {
continue
}
slog.InfoContext(ctx, "evicted download artifact (LRU)", "component", "downloads", "artifact_id", a.ID, "bytes", a.FileSize)
// The bytes are now guaranteed to be reclaimed — inline here, or by the
// orphan processor from the intent row — so the budget can be credited
// as soon as the row is gone.
freedInline := m.deleteEvictedArtifactBytes(ctx, a)
total -= a.FileSize
slog.InfoContext(ctx, "evicted download artifact (LRU)", "component", "downloads",
"artifact_id", a.ID, "bytes", a.FileSize, "bytes_freed_inline", freedInline)
}
return nil
}
// enqueueEvictionCleanup durably records how an artifact's bytes are to be
// removed once its row is deleted. Returns false when the intent could not be
// persisted, in which case the caller must leave the row in place.
func (m *ArtifactManager) enqueueEvictionCleanup(ctx context.Context, a *Artifact) bool {
if a.OriginArtifactID != "" {
if err := m.repo.EnqueueRemoteOrphan(ctx, a.ID, a.OriginNodeID, a.OriginNodeURL, a.OriginArtifactID); err != nil {
slog.WarnContext(ctx, "recording remote artifact cleanup intent failed", "component", "downloads", "artifact_id", a.ID, "node", a.OriginNodeURL, "error", err)
return false
}
return true
}
if a.OutputPath == "" {
return true
}
if err := m.repo.EnqueueLocalOrphan(ctx, a.ID, a.OutputPath); err != nil {
slog.WarnContext(ctx, "recording local artifact cleanup intent failed", "component", "downloads", "artifact_id", a.ID, "error", err)
return false
}
return true
}
// cleanupLocalOrphans retries filesystem deletions recorded before their
// artifact row was removed. Each candidate is verified against the artifact
// table first: paths carry the artifact id, so no successor row can ever claim
// the same path — a row that still exists means the eviction never happened and
// the file is live.
func (m *ArtifactManager) cleanupLocalOrphans(ctx context.Context) {
budget := m.remoteCleanupBudget
if budget <= 0 {
budget = defaultRemoteCleanupBudget
}
cleanupCtx, cancel := context.WithTimeout(ctx, budget)
defer cancel()
orphans, err := m.repo.ListLocalOrphansDue(cleanupCtx, 100)
if err != nil {
slog.WarnContext(ctx, "listing local artifact cleanup queue failed", "component", "downloads", "error", err)
return
}
for _, orphan := range orphans {
if cleanupCtx.Err() != nil {
return
}
_, err := m.repo.GetByID(cleanupCtx, orphan.DownloadArtifactID)
switch {
case err == nil:
// The row survived (eviction refused): the file is still owned.
if derr := m.repo.DeleteLocalOrphan(cleanupCtx, orphan.ID); derr != nil {
slog.WarnContext(ctx, "clearing local artifact cleanup row failed", "component", "downloads", "artifact_id", orphan.DownloadArtifactID, "error", derr)
}
continue
case !errors.Is(err, ErrNotFound):
slog.WarnContext(ctx, "verifying local artifact cleanup candidate failed", "component", "downloads", "artifact_id", orphan.DownloadArtifactID, "error", err)
continue
}
if !removeLocalArtifactFiles(ctx, orphan.DownloadArtifactID, orphan.OutputPath) {
if berr := m.repo.BumpLocalOrphanRetry(cleanupCtx, orphan.ID); berr != nil {
slog.WarnContext(ctx, "deferring local artifact cleanup row failed", "component", "downloads", "artifact_id", orphan.DownloadArtifactID, "error", berr)
}
continue
}
if derr := m.repo.DeleteLocalOrphan(cleanupCtx, orphan.ID); derr != nil {
slog.WarnContext(ctx, "clearing local artifact cleanup row failed", "component", "downloads", "artifact_id", orphan.DownloadArtifactID, "error", derr)
}
}
}
// removeLocalArtifactFiles unlinks an artifact's output file and its .part
// leftover, tolerating an already-missing file.
func removeLocalArtifactFiles(ctx context.Context, artifactID, outputPath string) bool {
if outputPath == "" {
return true
}
if err := os.Remove(outputPath); err != nil && !os.IsNotExist(err) {
slog.WarnContext(ctx, "removing artifact file failed", "component", "downloads", "artifact_id", artifactID, "error", err)
return false
}
if err := os.Remove(outputPath + ".part"); err != nil && !os.IsNotExist(err) {
slog.WarnContext(ctx, "removing artifact partial failed", "component", "downloads", "artifact_id", artifactID, "error", err)
return false
}
return true
}
func (m *ArtifactManager) cleanupRemoteOrphans(ctx context.Context) {
lifecycle, ok := m.preparer.(remoteArtifactLifecycle)
if !ok {
@@ -933,26 +1025,38 @@ func (m *ArtifactManager) sweepStale(ctx context.Context) {
// removeArtifact deletes an artifact's output file, its .part leftover, and
// its row. Used by the hygiene sweep for rows nothing can serve again.
func (m *ArtifactManager) removeArtifact(ctx context.Context, a *Artifact, reason string) {
if a.Status == ArtifactReady {
deleted, err := m.repo.DeleteReadyIfEvictable(ctx, a.ID)
if err != nil {
slog.WarnContext(ctx, "deleting swept artifact row failed", "component", "downloads", "artifact_id", a.ID, "error", err)
return
}
if !deleted {
return
}
_ = m.deleteEvictedArtifactBytes(ctx, a)
} else {
if !m.deleteArtifactBytes(ctx, a) {
return
}
if err := m.repo.DeleteArtifact(ctx, a.ID); err != nil {
slog.WarnContext(ctx, "deleting swept artifact row failed", "component", "downloads", "artifact_id", a.ID, "error", err)
return
}
// Both branches record the byte-cleanup intent before removing the row, so
// a failure after the delete defers to the orphan processor rather than
// leaking the file.
if !m.enqueueEvictionCleanup(ctx, a) {
return
}
slog.InfoContext(ctx, "swept stale download artifact", "component", "downloads", "artifact_id", a.ID, "reason", reason, "bytes", a.FileSize)
// A non-ready row is not necessarily dead: Ensure requeues a failed row and
// a download may have linked it since this candidate list was built, so both
// branches go through the active-link fence (plus the FK).
ready := a.Status == ArtifactReady
var deleted bool
var err error
if ready {
deleted, err = m.repo.DeleteReadyIfEvictable(ctx, a.ID)
} else {
deleted, err = m.repo.DeleteIfEvictable(ctx, a.ID)
}
if err != nil {
msg := "deleting swept artifact row failed"
if ready {
msg = "evicting swept ready artifact row failed"
}
slog.WarnContext(ctx, msg, "component", "downloads", "artifact_id", a.ID, "error", err)
return
}
if !deleted {
slog.InfoContext(ctx, "swept artifact retained; a download links it", "component", "downloads", "artifact_id", a.ID, "reason", reason)
return
}
freedInline := m.deleteEvictedArtifactBytes(ctx, a)
slog.InfoContext(ctx, "swept stale download artifact", "component", "downloads",
"artifact_id", a.ID, "reason", reason, "bytes", a.FileSize, "bytes_freed_inline", freedInline)
}
func (m *ArtifactManager) deleteArtifactBytes(ctx context.Context, a *Artifact) bool {
@@ -968,45 +1072,31 @@ func (m *ArtifactManager) deleteArtifactBytes(ctx context.Context, a *Artifact)
}
return true
}
if a.OutputPath == "" {
return true
}
if err := os.Remove(a.OutputPath); err != nil && !os.IsNotExist(err) {
slog.WarnContext(ctx, "removing artifact file failed", "component", "downloads", "artifact_id", a.ID, "error", err)
return false
}
if err := os.Remove(a.OutputPath + ".part"); err != nil && !os.IsNotExist(err) {
slog.WarnContext(ctx, "removing artifact partial failed", "component", "downloads", "artifact_id", a.ID, "error", err)
return false
}
return true
return removeLocalArtifactFiles(ctx, a.ID, a.OutputPath)
}
// deleteEvictedArtifactBytes removes bytes after the artifact row is already
// gone. Remote locators are unique per attempt, so they are always deleted
// (and re-queued for orphan cleanup on failure). Local output paths are
// deterministic: skip the unlink if a replacement job already claimed the path.
// deleteEvictedArtifactBytes removes bytes inline after the artifact row is
// already gone. A cleanup intent was persisted before the row delete, so this
// is only the fast path: on success the now-redundant intent row is cleared, on
// failure it is left for the orphan processor to retry. Paths carry the
// artifact id, so no replacement row can ever own this file.
func (m *ArtifactManager) deleteEvictedArtifactBytes(ctx context.Context, a *Artifact) bool {
if a.OriginArtifactID != "" {
if m.deleteArtifactBytes(ctx, a) {
return true
}
m.enqueueRemoteCleanup(ctx, a.ID, PreparedArtifact{
OriginNodeID: a.OriginNodeID, OriginNodeURL: a.OriginNodeURL, OriginArtifactID: a.OriginArtifactID,
}, false)
if !m.deleteArtifactBytes(ctx, a) {
return false
}
if a.OutputPath != "" {
inUse, err := m.repo.OutputPathInUse(ctx, a.OutputPath)
if err != nil {
slog.WarnContext(ctx, "artifact output path reuse check failed", "component", "downloads", "artifact_id", a.ID, "error", err)
return false
}
if inUse {
return true
if a.OriginArtifactID != "" {
// The remote queue dedupes on (origin_node_id, origin_artifact_id) and
// its processor re-checks ownership, so a stale row is harmless; clear
// it anyway to keep the queue short.
if err := m.repo.DeleteRemoteOrphanByLocator(ctx, a.OriginNodeID, a.OriginArtifactID); err != nil {
slog.WarnContext(ctx, "clearing remote artifact cleanup intent failed", "component", "downloads", "artifact_id", a.ID, "error", err)
}
return true
}
return m.deleteArtifactBytes(ctx, a)
if err := m.repo.DeleteLocalOrphanByPath(ctx, a.OutputPath); err != nil {
slog.WarnContext(ctx, "clearing local artifact cleanup intent failed", "component", "downloads", "artifact_id", a.ID, "error", err)
}
return true
}
// backoffFor returns the retry delay for the next attempt after a failure.
+103 -47
View File
@@ -9,6 +9,7 @@ import (
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgconn"
"github.com/jackc/pgx/v5/pgxpool"
)
@@ -16,19 +17,67 @@ const downloadColumns = `id, user_id, profile_id, device_id, media_file_id, cont
kind, status, format, quality, effective_quality, target_bitrate_kbps, revision, artifact_id, file_size, bytes_sent, error_message,
created_at, updated_at, completed_at`
const insertDownloadSQL = `INSERT INTO downloads (id, user_id, profile_id, device_id, media_file_id, content_id,
episode_id, batch_id, kind, status, format, quality, effective_quality, target_bitrate_kbps, revision, artifact_id, file_size, bytes_sent, error_message,
created_at, updated_at, completed_at)
// downloadInsertColumns is the number of columns in downloadColumns; the insert
// placeholder lists and the batch builder both derive their arity from it.
const downloadInsertColumns = 22
const insertDownloadSQL = `INSERT INTO downloads (` + downloadColumns + `)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22)`
// insertDownloadRequiringArtifactSQL refuses to create a row that would point
// at an artifact LRU/hygiene already deleted. $16 is artifact_id.
const insertDownloadRequiringArtifactSQL = `INSERT INTO downloads (id, user_id, profile_id, device_id, media_file_id, content_id,
episode_id, batch_id, kind, status, format, quality, effective_quality, target_bitrate_kbps, revision, artifact_id, file_size, bytes_sent, error_message,
created_at, updated_at, completed_at)
// insertDownloadRequiringArtifactSQL is a fast-path filter: it skips the insert
// when the artifact row is already gone, so the common eviction case surfaces as
// a clean zero-row result instead of a constraint error. It is NOT the
// enforcement mechanism — under READ COMMITTED its subquery reads an independent
// snapshot and takes no conflicting lock, so a concurrent eviction can still
// commit alongside it. downloads_artifact_id_fkey is what makes a dangling
// artifact_id impossible; a 23503 from it maps to ErrArtifactEvicted below.
// $16 is artifact_id.
const insertDownloadRequiringArtifactSQL = `INSERT INTO downloads (` + downloadColumns + `)
SELECT $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22
WHERE EXISTS (SELECT 1 FROM download_artifacts WHERE id = $16)`
// downloadsArtifactFKConstraint is the foreign key from downloads.artifact_id to
// download_artifacts.id. Violating it in either direction means the artifact was
// evicted concurrently.
const downloadsArtifactFKConstraint = "downloads_artifact_id_fkey"
// isForeignKeyViolation reports whether err is a Postgres foreign-key violation
// (SQLSTATE 23503) raised by the named constraint.
func isForeignKeyViolation(err error, constraint string) bool {
var pgErr *pgconn.PgError
if !errors.As(err, &pgErr) {
return false
}
return pgErr.Code == "23503" && pgErr.ConstraintName == constraint
}
// execer is satisfied by both *pgxpool.Pool and pgx.Tx, so the insert helper
// works inside and outside a transaction.
type execer interface {
Exec(ctx context.Context, sql string, args ...any) (pgconn.CommandTag, error)
}
// execInsertDownload runs the right insert for d (fenced when it links an
// artifact) and maps both the zero-row fast path and the FK violation to
// ErrArtifactEvicted.
func (r *Repository) execInsertDownload(ctx context.Context, q execer, d *Download) error {
query := insertDownloadSQL
if d.ArtifactID != "" {
query = insertDownloadRequiringArtifactSQL
}
tag, err := q.Exec(ctx, query, r.insertArgs(d)...)
if err != nil {
if isForeignKeyViolation(err, downloadsArtifactFKConstraint) {
return ErrArtifactEvicted
}
return err
}
if d.ArtifactID != "" && tag.RowsAffected() == 0 {
return ErrArtifactEvicted
}
return nil
}
// Repository provides CRUD operations for the downloads table.
type Repository struct {
pool *pgxpool.Pool
@@ -134,17 +183,12 @@ func (r *Repository) insertArgs(d *Download) []any {
// while the artifact row still exists, so LRU eviction cannot leave a ready
// download pointing at a deleted id.
func (r *Repository) Create(ctx context.Context, d *Download) error {
query := insertDownloadSQL
if d.ArtifactID != "" {
query = insertDownloadRequiringArtifactSQL
}
tag, err := r.pool.Exec(ctx, query, r.insertArgs(d)...)
if err != nil {
if err := r.execInsertDownload(ctx, r.pool, d); err != nil {
if errors.Is(err, ErrArtifactEvicted) {
return err
}
return fmt.Errorf("inserting download: %w", err)
}
if d.ArtifactID != "" && tag.RowsAffected() == 0 {
return ErrArtifactEvicted
}
return nil
}
@@ -157,17 +201,12 @@ func (r *Repository) CreateBatch(ctx context.Context, downloads []*Download) err
defer func() { _ = tx.Rollback(ctx) }()
for _, d := range downloads {
query := insertDownloadSQL
if d.ArtifactID != "" {
query = insertDownloadRequiringArtifactSQL
}
tag, err := tx.Exec(ctx, query, r.insertArgs(d)...)
if err != nil {
if err := r.execInsertDownload(ctx, tx, d); err != nil {
if errors.Is(err, ErrArtifactEvicted) {
return err
}
return fmt.Errorf("inserting batch download: %w", err)
}
if d.ArtifactID != "" && tag.RowsAffected() == 0 {
return ErrArtifactEvicted
}
}
if err := tx.Commit(ctx); err != nil {
@@ -274,9 +313,11 @@ func (r *Repository) Delete(ctx context.Context, id string, userID int) error {
}
// CancelByID sets a download to canceled if it is still queued or downloading.
// The artifact link is dropped with the transition: a terminal row must not pin
// an artifact against eviction (downloads_artifact_id_fkey is ON DELETE RESTRICT).
func (r *Repository) CancelByID(ctx context.Context, id string, userID int) error {
tag, err := r.pool.Exec(ctx,
`UPDATE downloads SET status = 'cancelled', updated_at = now()
`UPDATE downloads SET status = 'cancelled', artifact_id = NULL, updated_at = now()
WHERE id = $1 AND user_id = $2 AND status IN ('queued', 'downloading')`,
id, userID,
)
@@ -346,6 +387,11 @@ func (r *Repository) GetManagedEntry(ctx context.Context, userID int, profileID,
// Callers have already classified the item as new.
func (r *Repository) CreateManagedEntry(ctx context.Context, d *Download) (*Download, error) {
if err := r.Create(ctx, d); err != nil {
// An evicted artifact is retryable by the caller (re-Ensure, then
// re-insert); returning some unrelated existing row would hide it.
if errors.Is(err, ErrArtifactEvicted) {
return nil, err
}
if existing, gerr := r.GetManagedEntry(ctx, d.UserID, d.ProfileID, d.DeviceID, d.ContentID, d.EpisodeID); gerr == nil {
return existing, nil
}
@@ -411,11 +457,9 @@ func (r *Repository) CreateManagedEntriesBatch(ctx context.Context, ds []*Downlo
if len(ds) == 0 {
return nil, nil
}
const insertCols = 22
const insertCols = downloadInsertColumns
var sb strings.Builder
sb.WriteString(`INSERT INTO downloads (id, user_id, profile_id, device_id, media_file_id, content_id,
episode_id, batch_id, kind, status, format, quality, effective_quality, target_bitrate_kbps, revision, artifact_id, file_size, bytes_sent, error_message,
created_at, updated_at, completed_at) VALUES `)
sb.WriteString(`INSERT INTO downloads (` + downloadColumns + `) VALUES `)
args := make([]any, 0, len(ds)*insertCols)
for i, d := range ds {
if i > 0 {
@@ -435,14 +479,31 @@ func (r *Repository) CreateManagedEntriesBatch(ctx context.Context, ds []*Downlo
sb.WriteString(` ON CONFLICT (user_id, profile_id, device_id, content_id, (COALESCE(episode_id, ''))) WHERE device_id IS NOT NULL DO NOTHING RETURNING ` + downloadColumns)
rows, err := r.pool.Query(ctx, sb.String(), args...)
if err != nil {
// Managed originals carry no artifact today, so this is future-proofing:
// any artifact-backed row in the batch gets the same retryable error the
// single-row path returns.
if isForeignKeyViolation(err, downloadsArtifactFKConstraint) {
return nil, ErrArtifactEvicted
}
return nil, fmt.Errorf("batch inserting managed entries: %w", err)
}
defer rows.Close()
return scanDownloads(rows)
inserted, err := scanDownloads(rows)
if err != nil {
if isForeignKeyViolation(err, downloadsArtifactFKConstraint) {
return nil, ErrArtifactEvicted
}
return nil, err
}
return inserted, nil
}
// ReplaceManagedEntry updates an existing managed row to a new file/quality
// target, incrementing its revision so clients can replace stale local bytes.
// An artifact evicted concurrently surfaces as downloads_artifact_id_fkey
// (23503) and maps to ErrArtifactEvicted, so a zero-row result means only that
// the identity/revision fence rejected the write — another writer won the row,
// and the caller gets that winner back.
func (r *Repository) ReplaceManagedEntry(ctx context.Context, existing *Download, replacement *Download) (*Download, error) {
query := `UPDATE downloads SET
media_file_id = $6,
@@ -461,7 +522,6 @@ func (r *Repository) ReplaceManagedEntry(ctx context.Context, existing *Download
revision = revision + 1,
updated_at = now()
WHERE id = $1 AND user_id = $2 AND profile_id = $3 AND device_id = $4 AND revision = $5
AND ($14::text IS NULL OR EXISTS (SELECT 1 FROM download_artifacts WHERE id = $14))
RETURNING ` + downloadColumns
row := r.pool.QueryRow(ctx, query,
existing.ID, existing.UserID, existing.ProfileID, existing.DeviceID, existing.Revision,
@@ -472,20 +532,11 @@ func (r *Repository) ReplaceManagedEntry(ctx context.Context, existing *Download
d, err := scanDownload(row)
if err != nil {
if errors.Is(err, ErrNotFound) {
if replacement.ArtifactID != "" {
var artifactExists bool
if aerr := r.pool.QueryRow(ctx,
`SELECT EXISTS(SELECT 1 FROM download_artifacts WHERE id = $1)`,
replacement.ArtifactID,
).Scan(&artifactExists); aerr != nil {
return nil, fmt.Errorf("replacing managed download: %w", aerr)
}
if !artifactExists {
return nil, ErrArtifactEvicted
}
}
return r.GetManagedEntry(ctx, existing.UserID, existing.ProfileID, existing.DeviceID, existing.ContentID, existing.EpisodeID)
}
if isForeignKeyViolation(err, downloadsArtifactFKConstraint) {
return nil, ErrArtifactEvicted
}
return nil, fmt.Errorf("replacing managed download: %w", err)
}
return d, nil
@@ -649,10 +700,12 @@ func (r *Repository) MarkLinkedDownloadsReady(ctx context.Context, artifactID st
}
// MarkLinkedDownloadsFailed flips every preparing download linked to a failed
// artifact to failed, returning the affected rows for client notification.
// artifact to failed, returning the affected rows for client notification. The
// artifact link is dropped with the transition so the now-terminal row does not
// pin the artifact against eviction (the FK is ON DELETE RESTRICT).
func (r *Repository) MarkLinkedDownloadsFailed(ctx context.Context, artifactID, errMsg string) ([]*Download, error) {
rows, err := r.pool.Query(ctx,
`UPDATE downloads SET status = 'failed', error_message = $2, updated_at = now()
`UPDATE downloads SET status = 'failed', artifact_id = NULL, error_message = $2, updated_at = now()
WHERE artifact_id = $1 AND status = 'preparing'
RETURNING `+downloadColumns,
artifactID, errMsg,
@@ -692,8 +745,11 @@ func (r *Repository) ReconcileLinkedDownloads(ctx context.Context) (ready []*Dow
}
failedRows, err := r.pool.Query(ctx,
// Every SET expression reads the OLD row, so error_message still resolves
// against the pre-update artifact_id while the same statement clears it.
`UPDATE downloads SET status = 'failed',
error_message = COALESCE((SELECT NULLIF(a.error_message, '') FROM download_artifacts a WHERE a.id = downloads.artifact_id), 'artifact encode failed'),
artifact_id = NULL,
updated_at = now()
WHERE status = 'preparing' AND artifact_id IS NOT NULL
AND artifact_id IN (SELECT id FROM download_artifacts WHERE status = 'failed')
+81
View File
@@ -276,6 +276,10 @@ func TestReconcileLinkedDownloads(t *testing.T) {
}
arepo := NewArtifactRepository(f.pool)
t.Cleanup(func() {
// Cleanups run LIFO, so this fires before the fixture's downloads
// delete; unpin the links first or downloads_artifact_id_fkey
// (ON DELETE RESTRICT) refuses the artifact delete.
_, _ = f.pool.Exec(ctx, `UPDATE downloads SET artifact_id = NULL WHERE media_file_id = $1`, f.fileID)
_, _ = f.pool.Exec(ctx, `DELETE FROM download_artifacts WHERE media_file_id = $1`, f.fileID)
})
@@ -501,3 +505,80 @@ func TestRegisterManagedItemsBatchAndCount(t *testing.T) {
t.Fatalf("grown register = %d rows, want 2", len(grown))
}
}
// TestReplaceManagedEntryRevisionConflictReturnsWinner pins the simplified
// outcome contract: with the artifact foreign key doing the enforcement, a
// zero-row UPDATE can only mean the identity/revision fence rejected the write,
// so the caller gets the winning row rather than an eviction error.
func TestReplaceManagedEntryRevisionConflictReturnsWinner(t *testing.T) {
f := seedManagedFixture(t)
ctx := context.Background()
id := f.createManagedEntry(t)
existing, err := f.repo.GetByID(ctx, id)
if err != nil {
t.Fatalf("get managed entry: %v", err)
}
// A concurrent writer bumps the revision out from under the snapshot.
if _, err := f.pool.Exec(ctx, `UPDATE downloads SET revision = revision + 1 WHERE id = $1`, id); err != nil {
t.Fatalf("bump revision: %v", err)
}
replacement := &Download{
UserID: f.userID, ProfileID: f.profileA, DeviceID: f.deviceA,
MediaFileID: f.fileID, ContentID: f.contentID, Kind: KindQueued,
Status: StatusReady, Format: FormatOriginal, Quality: QualityOriginal,
EffectiveQuality: QualityOriginal, FileSize: 2048,
}
row, err := f.repo.ReplaceManagedEntry(ctx, existing, replacement)
if err != nil {
t.Fatalf("ReplaceManagedEntry on revision conflict: %v", err)
}
if row == nil || row.ID != id {
t.Fatalf("ReplaceManagedEntry returned %+v, want the winning row %q", row, id)
}
if row.Revision != existing.Revision+1 {
t.Fatalf("winning revision = %d, want %d (the concurrent writer's)", row.Revision, existing.Revision+1)
}
}
// TestReplaceManagedEntryEvictedArtifact pins the other half: an artifact
// evicted between Ensure and the replace surfaces as ErrArtifactEvicted (via
// downloads_artifact_id_fkey), never as a silent no-op.
func TestReplaceManagedEntryEvictedArtifact(t *testing.T) {
f := seedManagedFixture(t)
ctx := context.Background()
var present *string
if err := f.pool.QueryRow(ctx, `SELECT to_regclass('public.download_artifacts')::text`).Scan(&present); err != nil {
t.Fatalf("check download_artifacts: %v", err)
}
if present == nil {
t.Skip("download_artifacts migration has not been applied")
}
var fkPresent bool
if err := f.pool.QueryRow(ctx,
`SELECT EXISTS(SELECT 1 FROM pg_constraint WHERE conname = 'downloads_artifact_id_fkey')`,
).Scan(&fkPresent); err != nil {
t.Fatalf("check artifact fk: %v", err)
}
if !fkPresent {
t.Skip("downloads artifact foreign key migration has not been applied")
}
id := f.createManagedEntry(t)
existing, err := f.repo.GetByID(ctx, id)
if err != nil {
t.Fatalf("get managed entry: %v", err)
}
missingArtifactID := fmt.Sprintf("gone-artifact-%d", time.Now().UnixNano())
replacement := &Download{
UserID: f.userID, ProfileID: f.profileA, DeviceID: f.deviceA,
MediaFileID: f.fileID, ContentID: f.contentID, Kind: KindQueued,
Status: StatusReady, Format: FormatTranscode, Quality: QualityOriginal,
EffectiveQuality: QualityOriginal, ArtifactID: missingArtifactID, FileSize: 2048,
}
if _, err := f.repo.ReplaceManagedEntry(ctx, existing, replacement); !errors.Is(err, ErrArtifactEvicted) {
t.Fatalf("ReplaceManagedEntry with missing artifact = %v, want ErrArtifactEvicted", err)
}
}
+54 -38
View File
@@ -455,18 +455,27 @@ func (s *Service) createArtifactDownload(ctx context.Context, userID int, req Cr
const artifactLinkAttempts = 2
func (s *Service) replaceManagedArtifactDownload(ctx context.Context, existing *Download, userID int, req CreateRequest, file *models.MediaFile, decision QualityDecision) (*Download, error) {
// withEnsuredArtifact runs attempt against a freshly ensured artifact, retrying
// once when the artifact is evicted between Ensure and the linking write. Ensure
// must stay inside the loop: an evicted artifact has to be recreated before the
// next attempt can link anything.
func (s *Service) withEnsuredArtifact(
ctx context.Context,
file *models.MediaFile,
decision QualityDecision,
attempt func(artifact *Artifact, status string, size int64) (*Download, error),
) (*Download, error) {
var last error
for attempt := 0; attempt < artifactLinkAttempts; attempt++ {
for i := 0; i < artifactLinkAttempts; i++ {
artifact, err := s.artifacts.Ensure(ctx, file, decision.DeliveryFormat, decision.PrepareTarget)
if err != nil {
return nil, err
}
status, size := artifactRowStatus(artifact, file)
replacement := buildManagedDownload(userID, req.ProfileID, req.DeviceID, managedItem{file: file, contentID: file.ContentID, episodeID: file.EpisodeID}, decision, "", status, size, artifact.ID)
row, err := s.reuseOrReplaceManaged(ctx, existing, replacement)
row, err := attempt(artifact, status, size)
if errors.Is(err, ErrArtifactEvicted) {
last = err
slog.DebugContext(ctx, "retrying download after artifact eviction", "component", "downloads", "artifact_id", artifact.ID, "attempt", i+1)
continue
}
return row, err
@@ -474,21 +483,36 @@ func (s *Service) replaceManagedArtifactDownload(ctx context.Context, existing *
return nil, last
}
func (s *Service) replaceManagedArtifactDownload(ctx context.Context, existing *Download, userID int, req CreateRequest, file *models.MediaFile, decision QualityDecision) (*Download, error) {
return s.withEnsuredArtifact(ctx, file, decision, func(artifact *Artifact, status string, size int64) (*Download, error) {
row, err := s.replaceManagedWithArtifact(ctx, &existing, userID, req, file, decision, artifact, status, size)
return row, err
})
}
// replaceManagedWithArtifact points an existing managed entry at artifact. On
// eviction it refreshes the caller's *existing snapshot so the next attempt's
// revision fence can succeed against whatever the row looks like now.
func (s *Service) replaceManagedWithArtifact(ctx context.Context, existing **Download, userID int, req CreateRequest, file *models.MediaFile, decision QualityDecision, artifact *Artifact, status string, size int64) (*Download, error) {
replacement := buildManagedDownload(userID, req.ProfileID, req.DeviceID, managedItem{file: file, contentID: file.ContentID, episodeID: file.EpisodeID}, decision, "", status, size, artifact.ID)
row, err := s.reuseOrReplaceManaged(ctx, *existing, replacement)
if errors.Is(err, ErrArtifactEvicted) {
if fresh, gerr := s.repo.GetManagedEntry(ctx, userID, req.ProfileID, req.DeviceID, file.ContentID, file.EpisodeID); gerr == nil {
*existing = fresh
}
}
return row, err
}
func (s *Service) insertArtifactDownload(ctx context.Context, userID int, req CreateRequest, file *models.MediaFile, decision QualityDecision, managed bool) (*Download, error) {
var last error
for attempt := 0; attempt < artifactLinkAttempts; attempt++ {
artifact, err := s.artifacts.Ensure(ctx, file, decision.DeliveryFormat, decision.PrepareTarget)
if err != nil {
// The device upsert is idempotent and independent of the artifact, so it
// runs once rather than on every retry.
if managed {
if err := s.repo.EnsureDevice(ctx, userID, req.ProfileID, req.DeviceID, req.DeviceName, req.DevicePlatform); err != nil {
return nil, err
}
status, size := artifactRowStatus(artifact, file)
if managed {
if err := s.repo.EnsureDevice(ctx, userID, req.ProfileID, req.DeviceID, req.DeviceName, req.DevicePlatform); err != nil {
return nil, err
}
}
}
return s.withEnsuredArtifact(ctx, file, decision, func(artifact *Artifact, status string, size int64) (*Download, error) {
id, err := idgen.NextID()
if err != nil {
return nil, fmt.Errorf("generating download ID: %w", err)
@@ -516,30 +540,22 @@ func (s *Service) insertArtifactDownload(ctx context.Context, userID int, req Cr
d.ProfileID = req.ProfileID
d.DeviceID = req.DeviceID
}
if err := s.repo.Create(ctx, d); err != nil {
if errors.Is(err, ErrArtifactEvicted) {
last = err
continue
}
if managed {
if existing, gerr := s.repo.GetManagedEntry(ctx, userID, req.ProfileID, req.DeviceID, file.ContentID, file.EpisodeID); gerr == nil {
replacement := buildManagedDownload(userID, req.ProfileID, req.DeviceID, managedItem{file: file, contentID: file.ContentID, episodeID: file.EpisodeID}, decision, "", status, size, artifact.ID)
row, rerr := s.reuseOrReplaceManaged(ctx, existing, replacement)
if errors.Is(rerr, ErrArtifactEvicted) {
last = rerr
continue
}
if rerr != nil {
return nil, rerr
}
return row, nil
}
}
err = s.repo.Create(ctx, d)
if err == nil {
return d, nil
}
if errors.Is(err, ErrArtifactEvicted) {
return nil, err
}
return d, nil
}
return nil, last
// A concurrent create won this managed identity: fall back to replacing
// the winning row, reusing the artifact already ensured this attempt.
if managed {
if existing, gerr := s.repo.GetManagedEntry(ctx, userID, req.ProfileID, req.DeviceID, file.ContentID, file.EpisodeID); gerr == nil {
return s.replaceManagedWithArtifact(ctx, &existing, userID, req, file, decision, artifact, status, size)
}
}
return nil, err
})
}
// artifactRowStatus maps an ensured artifact to the download row status and
+5
View File
@@ -52,6 +52,10 @@ func TestArtifactQuotaCheckedBeforeEnqueue(t *testing.T) {
t.Fatalf("seed second media file: %v", err)
}
t.Cleanup(func() {
// Cleanups run LIFO, so this fires before the fixture's downloads
// delete; unpin the links first or downloads_artifact_id_fkey
// (ON DELETE RESTRICT) refuses the artifact delete.
_, _ = f.pool.Exec(ctx, `UPDATE downloads SET artifact_id = NULL WHERE media_file_id IN ($1, $2)`, f.fileID, fileID2)
_, _ = f.pool.Exec(ctx, `DELETE FROM download_artifacts WHERE media_file_id IN ($1, $2)`, f.fileID, fileID2)
_, _ = f.pool.Exec(ctx, `DELETE FROM media_files WHERE id = $1`, fileID2)
})
@@ -136,6 +140,7 @@ func TestConcurrentArtifactCreatesCannotBypassQuotaDB(t *testing.T) {
}
}
t.Cleanup(func() {
_, _ = f.pool.Exec(ctx, `UPDATE downloads SET artifact_id = NULL WHERE media_file_id = ANY($1)`, fileIDs)
_, _ = f.pool.Exec(ctx, `DELETE FROM download_artifacts WHERE media_file_id = ANY($1)`, fileIDs)
_, _ = f.pool.Exec(ctx, `DELETE FROM media_files WHERE id = ANY($1)`, fileIDs)
})
@@ -0,0 +1,25 @@
-- +goose Up
-- Terminal rows must not pin an artifact: the eviction fence treats their links
-- as dead, so an ON DELETE RESTRICT reference from one would block cleanup forever.
UPDATE public.downloads
SET artifact_id = NULL
WHERE artifact_id IS NOT NULL AND status IN ('cancelled', 'failed', 'revoked');
-- Heal legacy dangling references left by the pre-fence check-then-delete race.
UPDATE public.downloads
SET artifact_id = NULL
WHERE artifact_id IS NOT NULL
AND NOT EXISTS (SELECT 1 FROM public.download_artifacts a WHERE a.id = public.downloads.artifact_id);
-- The real referential-integrity check. Application-level EXISTS fences read
-- independent snapshots under READ COMMITTED and take no conflicting locks, so
-- a concurrent evict/link pair can both commit; the FK's RI trigger takes a
-- FOR KEY SHARE row lock and re-checks with a fresh snapshot after blocking,
-- which is what actually makes a dangling artifact_id impossible.
ALTER TABLE public.downloads
ADD CONSTRAINT downloads_artifact_id_fkey
FOREIGN KEY (artifact_id) REFERENCES public.download_artifacts(id) ON DELETE RESTRICT;
-- +goose Down
ALTER TABLE public.downloads
DROP CONSTRAINT IF EXISTS downloads_artifact_id_fkey;
@@ -0,0 +1,18 @@
-- +goose Up
CREATE TABLE public.download_artifact_local_orphans (
id bigserial PRIMARY KEY,
download_artifact_id text NOT NULL,
output_path text NOT NULL,
attempts integer NOT NULL DEFAULT 0,
next_retry_at timestamptz,
created_at timestamptz NOT NULL DEFAULT now(),
UNIQUE (output_path),
CONSTRAINT download_artifact_local_orphans_locator_check
CHECK (download_artifact_id <> '' AND output_path <> '')
);
CREATE INDEX download_artifact_local_orphans_due_idx
ON public.download_artifact_local_orphans (next_retry_at, created_at);
-- +goose Down
DROP TABLE IF EXISTS public.download_artifact_local_orphans;