diff --git a/internal/downloads/artifact.go b/internal/downloads/artifact.go index d0045eb6..a61ef0a1 100644 --- a/internal/downloads/artifact.go +++ b/internal/downloads/artifact.go @@ -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)) } diff --git a/internal/downloads/artifact_repo.go b/internal/downloads/artifact_repo.go index 674ebbe4..113e07cb 100644 --- a/internal/downloads/artifact_repo.go +++ b/internal/downloads/artifact_repo.go @@ -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) diff --git a/internal/downloads/artifact_repo_test.go b/internal/downloads/artifact_repo_test.go index 3aa39f9e..0d514346 100644 --- a/internal/downloads/artifact_repo_test.go +++ b/internal/downloads/artifact_repo_test.go @@ -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) + } + } +} diff --git a/internal/downloads/artifact_test.go b/internal/downloads/artifact_test.go index 9624a6af..a9df2116 100644 --- a/internal/downloads/artifact_test.go +++ b/internal/downloads/artifact_test.go @@ -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 { diff --git a/internal/downloads/artifacts.go b/internal/downloads/artifacts.go index c5a0ae4d..b5eafd49 100644 --- a/internal/downloads/artifacts.go +++ b/internal/downloads/artifacts.go @@ -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. diff --git a/internal/downloads/repo.go b/internal/downloads/repo.go index 692ccd83..126a6230 100644 --- a/internal/downloads/repo.go +++ b/internal/downloads/repo.go @@ -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') diff --git a/internal/downloads/repo_managed_test.go b/internal/downloads/repo_managed_test.go index e17059ed..5971938b 100644 --- a/internal/downloads/repo_managed_test.go +++ b/internal/downloads/repo_managed_test.go @@ -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) + } +} diff --git a/internal/downloads/service.go b/internal/downloads/service.go index 469d0940..99cc695a 100644 --- a/internal/downloads/service.go +++ b/internal/downloads/service.go @@ -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 diff --git a/internal/downloads/service_quota_test.go b/internal/downloads/service_quota_test.go index e0806b3b..172ce0d7 100644 --- a/internal/downloads/service_quota_test.go +++ b/internal/downloads/service_quota_test.go @@ -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) }) diff --git a/migrations/sql/20260813155652_add_downloads_artifact_fk.sql b/migrations/sql/20260813155652_add_downloads_artifact_fk.sql new file mode 100644 index 00000000..ec9fb920 --- /dev/null +++ b/migrations/sql/20260813155652_add_downloads_artifact_fk.sql @@ -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; diff --git a/migrations/sql/20260813155653_add_download_artifact_local_orphans.sql b/migrations/sql/20260813155653_add_download_artifact_local_orphans.sql new file mode 100644 index 00000000..7f9a536f --- /dev/null +++ b/migrations/sql/20260813155653_add_download_artifact_local_orphans.sql @@ -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;